Skip to main content

script/dom/stream/
writablestream.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, RefCell};
6use std::collections::VecDeque;
7use std::mem;
8use std::ptr::{self};
9use std::rc::Rc;
10
11use dom_struct::dom_struct;
12use js::context::JSContext;
13use js::conversions::ToJSValConvertible;
14use js::jsapi::{Heap, JSObject};
15use js::jsval::{JSVal, ObjectValue, UndefinedValue};
16use js::realm::CurrentRealm;
17use js::rust::{
18    HandleObject as SafeHandleObject, HandleValue as SafeHandleValue,
19    MutableHandleValue as SafeMutableHandleValue,
20};
21use rustc_hash::FxHashMap;
22use script_bindings::cell::DomRefCell;
23use script_bindings::codegen::GenericBindings::MessagePortBinding::MessagePortMethods;
24use script_bindings::reflector::{Reflector, reflect_dom_object_with_proto};
25use servo_base::id::{MessagePortId, MessagePortIndex};
26use servo_constellation_traits::MessagePortImpl;
27
28use crate::dom::bindings::codegen::Bindings::QueuingStrategyBinding::{
29    QueuingStrategy, QueuingStrategySize,
30};
31use crate::dom::bindings::codegen::Bindings::UnderlyingSinkBinding::UnderlyingSink;
32use crate::dom::bindings::codegen::Bindings::WritableStreamBinding::WritableStreamMethods;
33use crate::dom::bindings::conversions::ConversionResult;
34use crate::dom::bindings::error::{Error, Fallible};
35use crate::dom::bindings::reflector::DomGlobal;
36use crate::dom::bindings::root::{Dom, DomRoot, MutNullableDom};
37use crate::dom::bindings::structuredclone::StructuredData;
38use crate::dom::bindings::transferable::Transferable;
39use crate::dom::domexception::{DOMErrorName, DOMException};
40use crate::dom::globalscope::GlobalScope;
41use crate::dom::messageport::MessagePort;
42use crate::dom::promise::{Promise, RootedPromise, TracedPromise};
43use crate::dom::promisenativehandler::{Callback, PromiseNativeHandler};
44use crate::dom::readablestream::{ReadableStream, get_type_and_value_from_message};
45use crate::dom::stream::countqueuingstrategy::{extract_high_water_mark, extract_size_algorithm};
46use crate::dom::stream::writablestreamdefaultcontroller::{
47    UnderlyingSinkType, WritableStreamDefaultController,
48};
49use crate::dom::stream::writablestreamdefaultwriter::WritableStreamDefaultWriter;
50use crate::realms::enter_auto_realm;
51
52impl js::gc::Rootable for AbortAlgorithmFulfillmentHandler {}
53
54/// The fulfillment handler for the abort steps of
55/// <https://streams.spec.whatwg.org/#writable-stream-finish-erroring>
56#[derive(JSTraceable, MallocSizeOf)]
57#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
58struct AbortAlgorithmFulfillmentHandler {
59    stream: Dom<WritableStream>,
60    abort_request_promise: TracedPromise,
61}
62
63impl Callback for AbortAlgorithmFulfillmentHandler {
64    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
65        // Resolve abortRequest’s promise with undefined.
66        self.abort_request_promise.resolve_native(cx, &());
67
68        // Perform ! WritableStreamRejectCloseAndClosedPromiseIfNeeded(stream).
69        self.stream
70            .as_rooted()
71            .reject_close_and_closed_promise_if_needed(cx);
72    }
73}
74
75impl js::gc::Rootable for AbortAlgorithmRejectionHandler {}
76
77/// The rejection handler for the abort steps of
78/// <https://streams.spec.whatwg.org/#writable-stream-finish-erroring>
79#[derive(JSTraceable, MallocSizeOf)]
80#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
81struct AbortAlgorithmRejectionHandler {
82    stream: Dom<WritableStream>,
83    abort_request_promise: TracedPromise,
84}
85
86impl Callback for AbortAlgorithmRejectionHandler {
87    fn callback(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) {
88        // Reject abortRequest’s promise with reason.
89        self.abort_request_promise.reject_native(cx, &reason);
90
91        // Perform ! WritableStreamRejectCloseAndClosedPromiseIfNeeded(stream).
92        self.stream
93            .as_rooted()
94            .reject_close_and_closed_promise_if_needed(cx);
95    }
96}
97
98impl js::gc::Rootable for PendingAbortRequest {}
99
100/// <https://streams.spec.whatwg.org/#pending-abort-request>
101#[derive(JSTraceable, MallocSizeOf)]
102#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
103struct PendingAbortRequest {
104    /// <https://streams.spec.whatwg.org/#pending-abort-request-promise>
105    promise: TracedPromise,
106
107    /// <https://streams.spec.whatwg.org/#pending-abort-request-reason>
108    #[ignore_malloc_size_of = "mozjs"]
109    reason: Box<Heap<JSVal>>,
110
111    /// <https://streams.spec.whatwg.org/#pending-abort-request-was-already-erroring>
112    was_already_erroring: bool,
113}
114
115/// <https://streams.spec.whatwg.org/#writablestream-state>
116#[derive(Clone, Copy, Debug, Default, JSTraceable, MallocSizeOf)]
117pub(crate) enum WritableStreamState {
118    #[default]
119    Writable,
120    Closed,
121    Erroring,
122    Errored,
123}
124
125/// <https://streams.spec.whatwg.org/#ws-class>
126#[dom_struct]
127pub struct WritableStream {
128    reflector_: Reflector,
129
130    /// <https://streams.spec.whatwg.org/#writablestream-backpressure>
131    backpressure: Cell<bool>,
132
133    /// <https://streams.spec.whatwg.org/#writablestream-closerequest>
134    close_request: DomRefCell<Option<TracedPromise>>,
135
136    /// <https://streams.spec.whatwg.org/#writablestream-controller>
137    controller: MutNullableDom<WritableStreamDefaultController>,
138
139    /// <https://streams.spec.whatwg.org/#writablestream-detached>
140    detached: Cell<bool>,
141
142    /// <https://streams.spec.whatwg.org/#writablestream-inflightwriterequest>
143    in_flight_write_request: DomRefCell<Option<TracedPromise>>,
144
145    /// <https://streams.spec.whatwg.org/#writablestream-inflightcloserequest>
146    in_flight_close_request: DomRefCell<Option<TracedPromise>>,
147
148    /// <https://streams.spec.whatwg.org/#writablestream-pendingabortrequest>
149    pending_abort_request: DomRefCell<Option<PendingAbortRequest>>,
150
151    /// <https://streams.spec.whatwg.org/#writablestream-state>
152    state: Cell<WritableStreamState>,
153
154    /// <https://streams.spec.whatwg.org/#writablestream-storederror>
155    #[ignore_malloc_size_of = "mozjs"]
156    stored_error: Heap<JSVal>,
157
158    /// <https://streams.spec.whatwg.org/#writablestream-writer>
159    writer: MutNullableDom<WritableStreamDefaultWriter>,
160
161    /// <https://streams.spec.whatwg.org/#writablestream-writerequests>
162    write_requests: DomRefCell<VecDeque<TracedPromise>>,
163}
164
165impl WritableStream {
166    /// <https://streams.spec.whatwg.org/#initialize-writable-stream>
167    fn new_inherited() -> WritableStream {
168        WritableStream {
169            reflector_: Reflector::new(),
170            backpressure: Default::default(),
171            close_request: Default::default(),
172            controller: Default::default(),
173            detached: Default::default(),
174            in_flight_write_request: Default::default(),
175            in_flight_close_request: Default::default(),
176            pending_abort_request: Default::default(),
177            state: Default::default(),
178            stored_error: Default::default(),
179            writer: Default::default(),
180            write_requests: Default::default(),
181        }
182    }
183
184    pub(crate) fn new_with_proto(
185        cx: &mut JSContext,
186        global: &GlobalScope,
187        proto: Option<SafeHandleObject>,
188    ) -> DomRoot<WritableStream> {
189        reflect_dom_object_with_proto(cx, Box::new(WritableStream::new_inherited()), global, proto)
190    }
191
192    /// Used as part of
193    /// <https://streams.spec.whatwg.org/#set-up-writable-stream-default-controller>
194    pub(crate) fn assert_no_controller(&self) {
195        assert!(self.controller.get().is_none());
196    }
197
198    /// Used as part of
199    /// <https://streams.spec.whatwg.org/#set-up-writable-stream-default-controller>
200    pub(crate) fn set_default_controller(&self, controller: &WritableStreamDefaultController) {
201        self.controller.set(Some(controller));
202    }
203
204    pub(crate) fn get_default_controller(&self) -> DomRoot<WritableStreamDefaultController> {
205        self.controller.get().expect("Controller should be set.")
206    }
207
208    pub(crate) fn is_writable(&self) -> bool {
209        matches!(self.state.get(), WritableStreamState::Writable)
210    }
211
212    pub(crate) fn is_erroring(&self) -> bool {
213        matches!(self.state.get(), WritableStreamState::Erroring)
214    }
215
216    pub(crate) fn is_errored(&self) -> bool {
217        matches!(self.state.get(), WritableStreamState::Errored)
218    }
219
220    pub(crate) fn is_closed(&self) -> bool {
221        matches!(self.state.get(), WritableStreamState::Closed)
222    }
223
224    pub(crate) fn has_in_flight_write_request(&self) -> bool {
225        self.in_flight_write_request.borrow().is_some()
226    }
227
228    /// <https://streams.spec.whatwg.org/#writable-stream-has-operation-marked-in-flight>
229    pub(crate) fn has_operations_marked_inflight(&self) -> bool {
230        let in_flight_write_requested = self.in_flight_write_request.borrow().is_some();
231        let in_flight_close_requested = self.in_flight_close_request.borrow().is_some();
232
233        in_flight_write_requested || in_flight_close_requested
234    }
235
236    /// <https://streams.spec.whatwg.org/#writablestream-storederror>
237    pub(crate) fn get_stored_error(&self, mut handle_mut: SafeMutableHandleValue) {
238        handle_mut.set(self.stored_error.get());
239    }
240
241    /// <https://streams.spec.whatwg.org/#writable-stream-finish-erroring>
242    pub(crate) fn finish_erroring(&self, cx: &mut JSContext, global: &GlobalScope) {
243        // Assert: stream.[[state]] is "erroring".
244        assert!(self.is_erroring());
245
246        // Assert: ! WritableStreamHasOperationMarkedInFlight(stream) is false.
247        assert!(!self.has_operations_marked_inflight());
248
249        // Set stream.[[state]] to "errored".
250        self.state.set(WritableStreamState::Errored);
251
252        // Perform ! stream.[[controller]].[[ErrorSteps]]().
253        let Some(controller) = self.controller.get() else {
254            unreachable!("Stream should have a controller.");
255        };
256        controller.perform_error_steps();
257
258        // Let storedError be stream.[[storedError]].
259        rooted!(&in(cx) let mut stored_error = UndefinedValue());
260        self.get_stored_error(stored_error.handle_mut());
261
262        // For each writeRequest of stream.[[writeRequests]]:
263        rooted!(&in(cx) let write_requests = mem::take(&mut *self.write_requests.borrow_mut()));
264        for request in write_requests.iter() {
265            // Reject writeRequest with storedError.
266            request.reject(cx, stored_error.handle());
267        }
268
269        // Set stream.[[writeRequests]] to an empty list.
270        // Done above with `drain`.
271
272        // If stream.[[pendingAbortRequest]] is undefined,
273        if self.pending_abort_request.borrow().is_none() {
274            // Perform ! WritableStreamRejectCloseAndClosedPromiseIfNeeded(stream).
275            self.reject_close_and_closed_promise_if_needed(cx);
276
277            // Return.
278            return;
279        }
280
281        // Let abortRequest be stream.[[pendingAbortRequest]].
282        // Set stream.[[pendingAbortRequest]] to undefined.
283        rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
284        if let Some(pending_abort_request) = &*pending_abort_request {
285            // If abortRequest’s was already erroring is true,
286            if pending_abort_request.was_already_erroring {
287                // Reject abortRequest’s promise with storedError.
288                pending_abort_request
289                    .promise
290                    .reject(cx, stored_error.handle());
291
292                // Perform ! WritableStreamRejectCloseAndClosedPromiseIfNeeded(stream).
293                self.reject_close_and_closed_promise_if_needed(cx);
294
295                // Return.
296                return;
297            }
298
299            // Let promise be ! stream.[[controller]].[[AbortSteps]](abortRequest’s reason).
300            rooted!(&in(cx) let mut reason = UndefinedValue());
301            reason.set(pending_abort_request.reason.get());
302            let promise = controller.abort_steps(cx, global, reason.handle());
303
304            // Upon fulfillment of promise,
305            rooted!(&in(cx) let mut fulfillment_handler = Some(AbortAlgorithmFulfillmentHandler {
306                stream: Dom::from_ref(self),
307                abort_request_promise: pending_abort_request.promise.clone(),
308            }));
309
310            // Upon rejection of promise with reason r,
311            rooted!(&in(cx) let mut rejection_handler = Some(AbortAlgorithmRejectionHandler {
312                stream: Dom::from_ref(self),
313                abort_request_promise: pending_abort_request.promise.clone(),
314            }));
315
316            let handler = PromiseNativeHandler::new(
317                cx,
318                global,
319                fulfillment_handler.take().map(|h| Box::new(h) as Box<_>),
320                rejection_handler.take().map(|h| Box::new(h) as Box<_>),
321            );
322
323            let mut realm = enter_auto_realm(cx, global);
324            let cx = &mut realm.current_realm();
325            promise.append_native_handler(cx, &handler);
326        }
327    }
328
329    /// <https://streams.spec.whatwg.org/#writable-stream-reject-close-and-closed-promise-if-needed>
330    fn reject_close_and_closed_promise_if_needed(&self, cx: &mut JSContext) {
331        // Assert: stream.[[state]] is "errored".
332        assert!(self.is_errored());
333
334        rooted!(&in(cx) let mut stored_error = UndefinedValue());
335        self.get_stored_error(stored_error.handle_mut());
336
337        // If stream.[[closeRequest]] is not undefined
338        rooted!(&in(cx) let close_request = self.close_request.borrow_mut().take());
339        if let Some(ref close_request) = *close_request {
340            // Assert: stream.[[inFlightCloseRequest]] is undefined.
341            assert!(self.in_flight_close_request.borrow().is_none());
342
343            // Reject stream.[[closeRequest]] with stream.[[storedError]].
344            close_request.reject_native(cx, &stored_error.handle())
345
346            // Set stream.[[closeRequest]] to undefined.
347            // Done with `take` above.
348        }
349
350        // Let writer be stream.[[writer]].
351        // If writer is not undefined,
352        if let Some(writer) = self.writer.get() {
353            // Reject writer.[[closedPromise]] with stream.[[storedError]].
354            writer.reject_closed_promise_with_stored_error(cx, &stored_error.handle());
355
356            // Set writer.[[closedPromise]].[[PromiseIsHandled]] to true.
357            writer.set_close_promise_is_handled(cx);
358        }
359    }
360
361    /// <https://streams.spec.whatwg.org/#writable-stream-close-queued-or-in-flight>
362    pub(crate) fn close_queued_or_in_flight(&self) -> bool {
363        let close_requested = self.close_request.borrow().is_some();
364        let in_flight_close_requested = self.in_flight_close_request.borrow().is_some();
365
366        close_requested || in_flight_close_requested
367    }
368
369    /// <https://streams.spec.whatwg.org/#writable-stream-finish-in-flight-write>
370    pub(crate) fn finish_in_flight_write(&self, cx: &mut JSContext) {
371        rooted!(&in(cx) let in_flight_write_request = self.in_flight_write_request.borrow_mut().take());
372        let Some(ref in_flight_write_request) = *in_flight_write_request else {
373            // Assert: stream.[[inFlightWriteRequest]] is not undefined.
374            unreachable!("Stream should have a write request");
375        };
376
377        // Resolve stream.[[inFlightWriteRequest]] with undefined.
378        in_flight_write_request.resolve_native(cx, &());
379
380        // Set stream.[[inFlightWriteRequest]] to undefined.
381        // Done above with `take`.
382    }
383
384    /// <https://streams.spec.whatwg.org/#writable-stream-start-erroring>
385    pub(crate) fn start_erroring(
386        &self,
387        cx: &mut JSContext,
388        global: &GlobalScope,
389        error: SafeHandleValue,
390    ) {
391        // Assert: stream.[[storedError]] is undefined.
392        assert!(self.stored_error.get().is_undefined());
393
394        // Assert: stream.[[state]] is "writable".
395        assert!(self.is_writable());
396
397        // Let controller be stream.[[controller]].
398        let Some(controller) = self.controller.get() else {
399            // Assert: controller is not undefined.
400            unreachable!("Stream should have a controller.");
401        };
402
403        // Set stream.[[state]] to "erroring".
404        self.state.set(WritableStreamState::Erroring);
405
406        // Set stream.[[storedError]] to reason.
407        self.stored_error.set(*error);
408
409        // Let writer be stream.[[writer]].
410        if let Some(writer) = self.writer.get() {
411            // If writer is not undefined, perform ! WritableStreamDefaultWriterEnsureReadyPromiseRejected
412            writer.ensure_ready_promise_rejected(cx, global, error);
413        }
414
415        // If ! WritableStreamHasOperationMarkedInFlight(stream) is false and controller.[[started]] is true
416        if !self.has_operations_marked_inflight() && controller.started() {
417            // perform ! WritableStreamFinishErroring
418            self.finish_erroring(cx, global);
419        }
420    }
421
422    /// <https://streams.spec.whatwg.org/#writable-stream-deal-with-rejection>
423    pub(crate) fn deal_with_rejection(
424        &self,
425        cx: &mut JSContext,
426        global: &GlobalScope,
427        error: SafeHandleValue,
428    ) {
429        // Let state be stream.[[state]].
430
431        // If state is "writable",
432        if self.is_writable() {
433            // Perform ! WritableStreamStartErroring(stream, error).
434            self.start_erroring(cx, global, error);
435
436            // Return.
437            return;
438        }
439
440        // Assert: state is "erroring".
441        assert!(self.is_erroring());
442
443        // Perform ! WritableStreamFinishErroring(stream).
444        self.finish_erroring(cx, global);
445    }
446
447    /// <https://streams.spec.whatwg.org/#writable-stream-mark-first-write-request-in-flight>
448    pub(crate) fn mark_first_write_request_in_flight(&self) {
449        let mut in_flight_write_request = self.in_flight_write_request.borrow_mut();
450        let mut write_requests = self.write_requests.borrow_mut();
451
452        // Assert: stream.[[inFlightWriteRequest]] is undefined.
453        assert!(in_flight_write_request.is_none());
454
455        // Assert: stream.[[writeRequests]] is not empty.
456        assert!(!write_requests.is_empty());
457
458        // Let writeRequest be stream.[[writeRequests]][0].
459        // Remove writeRequest from stream.[[writeRequests]].
460        // Set stream.[[inFlightWriteRequest]] to writeRequest.
461        *in_flight_write_request = write_requests.pop_front();
462    }
463
464    /// <https://streams.spec.whatwg.org/#writable-stream-mark-close-request-in-flight>
465    pub(crate) fn mark_close_request_in_flight(&self) {
466        let mut in_flight_close_request = self.in_flight_close_request.borrow_mut();
467        let mut close_request = self.close_request.borrow_mut();
468
469        // Assert: stream.[[inFlightCloseRequest]] is undefined.
470        assert!(in_flight_close_request.is_none());
471
472        // Assert: stream.[[closeRequest]] is not undefined.
473        assert!(close_request.is_some());
474
475        // Let closeRequest be stream.[[closeRequest]].
476        // Set stream.[[closeRequest]] to undefined.
477        // Set stream.[[inFlightCloseRequest]] to closeRequest.
478        *in_flight_close_request = close_request.take();
479    }
480
481    /// <https://streams.spec.whatwg.org/#writable-stream-finish-in-flight-close>
482    pub(crate) fn finish_in_flight_close(&self, cx: &mut JSContext) {
483        rooted!(&in(cx) let in_flight_close_request = self.in_flight_close_request.borrow_mut().take());
484        let Some(ref in_flight_close_request) = *in_flight_close_request else {
485            // Assert: stream.[[inFlightCloseRequest]] is not undefined.
486            unreachable!("in_flight_close_request must be Some");
487        };
488
489        // Resolve stream.[[inFlightCloseRequest]] with undefined.
490        in_flight_close_request.resolve_native(cx, &());
491
492        // Set stream.[[inFlightCloseRequest]] to undefined.
493        // Done with take above.
494
495        // Assert: stream.[[state]] is "writable" or "erroring".
496        assert!(self.is_writable() || self.is_erroring());
497
498        // If state is "erroring",
499        if self.is_erroring() {
500            // Set stream.[[storedError]] to undefined.
501            self.stored_error.set(UndefinedValue());
502
503            // If stream.[[pendingAbortRequest]] is not undefined,
504            rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
505            if let Some(pending_abort_request) = &*pending_abort_request {
506                // Resolve stream.[[pendingAbortRequest]]'s promise with undefined.
507                pending_abort_request.promise.resolve_native(cx, &());
508
509                // Set stream.[[pendingAbortRequest]] to undefined.
510                // Done above with `take`.
511            }
512        }
513
514        // Set stream.[[state]] to "closed".
515        self.state.set(WritableStreamState::Closed);
516
517        // Let writer be stream.[[writer]].
518        if let Some(writer) = self.writer.get() {
519            // If writer is not undefined,
520            // resolve writer.[[closedPromise]] with undefined.
521            writer.resolve_closed_promise_with_undefined(cx);
522        }
523
524        // Assert: stream.[[pendingAbortRequest]] is undefined.
525        assert!(self.pending_abort_request.borrow().is_none());
526
527        // Assert: stream.[[storedError]] is undefined.
528        assert!(self.stored_error.get().is_undefined());
529    }
530
531    /// <https://streams.spec.whatwg.org/#writable-stream-finish-in-flight-close-with-error>
532    pub(crate) fn finish_in_flight_close_with_error(
533        &self,
534        cx: &mut JSContext,
535        global: &GlobalScope,
536        error: SafeHandleValue,
537    ) {
538        rooted!(&in(cx) let in_flight_close_request = self.in_flight_close_request.borrow_mut().take());
539        let Some(ref in_flight_close_request) = *in_flight_close_request else {
540            // Assert: stream.[[inFlightCloseRequest]] is not undefined.
541            unreachable!("Inflight close request must be defined.");
542        };
543
544        // Reject stream.[[inFlightCloseRequest]] with error.
545        in_flight_close_request.reject_native(cx, &error);
546
547        // Set stream.[[inFlightCloseRequest]] to undefined.
548        // Done above with `take`.
549
550        // Assert: stream.[[state]] is "writable" or "erroring".
551        assert!(self.is_erroring() || self.is_writable());
552
553        // If stream.[[pendingAbortRequest]] is not undefined,
554        rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
555        if let Some(pending_abort_request) = &*pending_abort_request {
556            // Reject stream.[[pendingAbortRequest]]'s promise with error.
557            pending_abort_request.promise.reject_native(cx, &error);
558
559            // Set stream.[[pendingAbortRequest]] to undefined.
560            // Done above with `take`.
561        }
562
563        // Perform ! WritableStreamDealWithRejection(stream, error).
564        self.deal_with_rejection(cx, global, error);
565    }
566
567    /// <https://streams.spec.whatwg.org/#writable-stream-finish-in-flight-write-with-error>
568    pub(crate) fn finish_in_flight_write_with_error(
569        &self,
570        cx: &mut JSContext,
571        global: &GlobalScope,
572        error: SafeHandleValue,
573    ) {
574        rooted!(&in(cx) let in_flight_write_request = self.in_flight_write_request.borrow_mut().take());
575        let Some(ref in_flight_write_request) = *in_flight_write_request else {
576            // Assert: stream.[[inFlightWriteRequest]] is not undefined.
577            unreachable!("Inflight write request must be defined.");
578        };
579
580        // Reject stream.[[inFlightWriteRequest]] with error.
581        in_flight_write_request.reject_native(cx, &error);
582
583        // Set stream.[[inFlightWriteRequest]] to undefined.
584        // Done above with `take`.
585
586        // Assert: stream.[[state]] is "writable" or "erroring".
587        assert!(self.is_erroring() || self.is_writable());
588
589        // Perform ! WritableStreamDealWithRejection(stream, error).
590        self.deal_with_rejection(cx, global, error);
591    }
592
593    pub(crate) fn get_writer(&self) -> Option<DomRoot<WritableStreamDefaultWriter>> {
594        self.writer.get()
595    }
596
597    pub(crate) fn set_writer(&self, writer: Option<&WritableStreamDefaultWriter>) {
598        self.writer.set(writer);
599    }
600
601    pub(crate) fn set_backpressure(&self, backpressure: bool) {
602        self.backpressure.set(backpressure);
603    }
604
605    pub(crate) fn get_backpressure(&self) -> bool {
606        self.backpressure.get()
607    }
608
609    /// <https://streams.spec.whatwg.org/#is-writable-stream-locked>
610    pub(crate) fn is_locked(&self) -> bool {
611        // If stream.[[writer]] is undefined, return false.
612        // Return true.
613        self.get_writer().is_some()
614    }
615
616    /// <https://streams.spec.whatwg.org/#writable-stream-add-write-request>
617    pub(crate) fn add_write_request(
618        &self,
619        cx: &mut JSContext,
620        global: &GlobalScope,
621    ) -> RootedPromise {
622        // Assert: ! IsWritableStreamLocked(stream) is true.
623        assert!(self.is_locked());
624
625        // Assert: stream.[[state]] is "writable".
626        assert!(self.is_writable());
627
628        // Let promise be a new promise.
629        let promise = Promise::new_rooted(cx, global);
630
631        // Append promise to stream.[[writeRequests]].
632        self.write_requests
633            .borrow_mut()
634            .push_back(promise.to_traced());
635
636        // Return promise.
637        promise
638    }
639
640    // Returns the rooted controller of the stream, if any.
641    pub(crate) fn get_controller(&self) -> Option<DomRoot<WritableStreamDefaultController>> {
642        self.controller.get()
643    }
644
645    /// <https://streams.spec.whatwg.org/#writable-stream-abort>
646    pub(crate) fn abort(
647        &self,
648        cx: &mut CurrentRealm,
649        global: &GlobalScope,
650        provided_reason: SafeHandleValue,
651    ) -> RootedPromise {
652        // If stream.[[state]] is "closed" or "errored",
653        if self.is_closed() || self.is_errored() {
654            // return a promise resolved with undefined.
655            return Promise::new_resolved_rooted(cx, global, ());
656        }
657
658        // Signal abort on stream.[[controller]].[[abortController]] with reason.
659        self.get_controller()
660            .expect("Stream must have a controller.")
661            .signal_abort(cx, provided_reason);
662
663        // Let state be stream.[[state]].
664        let state = self.state.get();
665
666        // If state is "closed" or "errored", return a promise resolved with undefined.
667        if matches!(
668            state,
669            WritableStreamState::Closed | WritableStreamState::Errored
670        ) {
671            return Promise::new_resolved_rooted(cx, global, ());
672        }
673
674        // If stream.[[pendingAbortRequest]] is not undefined,
675        if self.pending_abort_request.borrow().is_some() {
676            // return stream.[[pendingAbortRequest]]'s promise.
677            return self
678                .pending_abort_request
679                .borrow()
680                .as_ref()
681                .expect("Pending abort request must be Some.")
682                .promise
683                .root();
684        }
685
686        // Assert: state is "writable" or "erroring".
687        assert!(self.is_writable() || self.is_erroring());
688
689        // Let wasAlreadyErroring be false.
690        let mut was_already_erroring = false;
691        rooted!(&in(cx) let undefined_reason = UndefinedValue());
692
693        // If state is "erroring",
694        let reason = if self.is_erroring() {
695            // Set wasAlreadyErroring to true.
696            was_already_erroring = true;
697
698            // Set reason to undefined.
699            undefined_reason.handle()
700        } else {
701            // Use the provided reason.
702            provided_reason
703        };
704
705        // Let promise be a new promise.
706        let promise = Promise::new_rooted(cx, global);
707
708        // Set stream.[[pendingAbortRequest]] to a new pending abort request
709        // whose promise is promise,
710        // reason is reason,
711        // and was already erroring is wasAlreadyErroring.
712        *self.pending_abort_request.borrow_mut() = Some(PendingAbortRequest {
713            promise: promise.to_traced(),
714            reason: Heap::boxed(reason.get()),
715            was_already_erroring,
716        });
717
718        // If wasAlreadyErroring is false,
719        if !was_already_erroring {
720            // perform ! WritableStreamStartErroring(stream, reason)
721            self.start_erroring(cx, global, reason);
722        }
723
724        // Return promise.
725        promise
726    }
727
728    /// <https://streams.spec.whatwg.org/#writable-stream-close>
729    pub(crate) fn close(&self, cx: &mut JSContext, global: &GlobalScope) -> RootedPromise {
730        // Let state be stream.[[state]].
731        // If state is "closed" or "errored",
732        if self.is_closed() || self.is_errored() {
733            // return a promise rejected with a TypeError exception.
734            let promise = Promise::new_rooted(cx, global);
735            promise.reject_error(cx, Error::Type(c"Stream is closed or errored.".to_owned()));
736            return promise;
737        }
738
739        // Assert: state is "writable" or "erroring".
740        assert!(self.is_writable() || self.is_erroring());
741
742        // Assert: ! WritableStreamCloseQueuedOrInFlight(stream) is false.
743        assert!(!self.close_queued_or_in_flight());
744
745        // Let promise be a new promise.
746        let promise = Promise::new_rooted(cx, global);
747
748        // Set stream.[[closeRequest]] to promise.
749        *self.close_request.borrow_mut() = Some(promise.to_traced());
750
751        // Let writer be stream.[[writer]].
752        // If writer is not undefined,
753        if let Some(writer) = self.writer.get() {
754            // and stream.[[backpressure]] is true,
755            // and state is "writable",
756            if self.get_backpressure() && self.is_writable() {
757                // resolve writer.[[readyPromise]] with undefined.
758                writer.resolve_ready_promise_with_undefined(cx);
759            }
760        }
761
762        // Perform ! WritableStreamDefaultControllerClose(stream.[[controller]]).
763        let Some(controller) = self.controller.get() else {
764            unreachable!("Stream must have a controller.");
765        };
766        controller.close(cx, global);
767
768        // Return promise.
769        promise
770    }
771
772    /// <https://streams.spec.whatwg.org/#writable-stream-default-writer-get-desired-size>
773    /// Note: implement as a stream method, as opposed to a writer one, for convenience.
774    pub(crate) fn get_desired_size(&self) -> Option<f64> {
775        // Let stream be writer.[[stream]].
776        // Stream is `self`.
777
778        // Let state be stream.[[state]].
779        // If state is "errored" or "erroring", return null.
780        if self.is_errored() || self.is_erroring() {
781            return None;
782        }
783
784        // If state is "closed", return 0.
785        if self.is_closed() {
786            return Some(0.);
787        }
788
789        let Some(controller) = self.controller.get() else {
790            unreachable!("Stream must have a controller.");
791        };
792        Some(controller.get_desired_size())
793    }
794
795    /// <https://streams.spec.whatwg.org/#acquire-writable-stream-default-writer>
796    pub(crate) fn aquire_default_writer(
797        &self,
798        cx: &mut CurrentRealm,
799        global: &GlobalScope,
800    ) -> Result<DomRoot<WritableStreamDefaultWriter>, Error> {
801        // Let writer be a new WritableStreamDefaultWriter object.
802        let writer = WritableStreamDefaultWriter::new(cx, global, None);
803
804        // Perform ? SetUpWritableStreamDefaultWriter(writer, stream).
805        writer.setup(cx, self)?;
806
807        // Return writer.
808        Ok(writer)
809    }
810
811    /// <https://streams.spec.whatwg.org/#writable-stream-update-backpressure>
812    pub(crate) fn update_backpressure(
813        &self,
814        cx: &mut JSContext,
815        backpressure: bool,
816        global: &GlobalScope,
817    ) {
818        // Assert: stream.[[state]] is "writable".
819        self.is_writable();
820
821        // Assert: ! WritableStreamCloseQueuedOrInFlight(stream) is false.
822        assert!(!self.close_queued_or_in_flight());
823
824        // Let writer be stream.[[writer]].
825        let writer = self.get_writer();
826
827        if let Some(writer) = writer {
828            // If writer is not undefined
829            if backpressure != self.get_backpressure() {
830                // and backpressure is not stream.[[backpressure]],
831                if backpressure {
832                    // If backpressure is true, set writer.[[readyPromise]] to a new promise.
833                    let promise = Promise::new(cx, global);
834                    writer.set_ready_promise(promise);
835                } else {
836                    // Otherwise,
837                    // Assert: backpressure is false.
838                    assert!(!backpressure);
839                    // Resolve writer.[[readyPromise]] with undefined.
840                    writer.resolve_ready_promise_with_undefined(cx);
841                }
842            }
843        }
844
845        // Set stream.[[backpressure]] to backpressure.
846        self.set_backpressure(backpressure);
847    }
848
849    /// <https://streams.spec.whatwg.org/#abstract-opdef-setupcrossrealmtransformwritable>
850    pub(crate) fn setup_cross_realm_transform_writable(
851        &self,
852        cx: &mut JSContext,
853        port: &MessagePort,
854    ) {
855        let port_id = port.message_port_id();
856        let global = self.global();
857
858        // Perform ! InitializeWritableStream(stream).
859        // Done in `new_inherited`.
860
861        // Let sizeAlgorithm be an algorithm that returns 1.
862        // Re-ordered because of the need to pass it to `new`.
863        let size_algorithm = extract_size_algorithm(cx, &QueuingStrategy::default());
864
865        // Note: other algorithms defined in the controller at call site.
866
867        // Let backpressurePromise be a new promise.
868        let backpressure_promise = Rc::new(RefCell::new(Some(Promise::new(cx, &global))));
869
870        // Let controller be a new WritableStreamDefaultController.
871        let controller = WritableStreamDefaultController::new(
872            cx,
873            &global,
874            UnderlyingSinkType::Transfer {
875                backpressure_promise: backpressure_promise.clone(),
876                port: Dom::from_ref(port),
877            },
878            1.0,
879            size_algorithm,
880        );
881
882        // Add a handler for port’s message event with the following steps:
883        // Add a handler for port’s messageerror event with the following steps:
884        rooted!(&in(cx) let cross_realm_transform_writable = CrossRealmTransformWritable {
885            controller: Dom::from_ref(&controller),
886            backpressure_promise,
887        });
888        global.note_cross_realm_transform_writable(&cross_realm_transform_writable, port_id);
889
890        // Enable port’s port message queue.
891        port.Start(cx);
892
893        // Perform ! SetUpWritableStreamDefaultController
894        controller
895            .setup(cx, &global, self)
896            .expect("Setup for transfer cannot fail");
897    }
898    /// <https://streams.spec.whatwg.org/#set-up-writable-stream-default-controller-from-underlying-sink>
899    #[allow(clippy::too_many_arguments)]
900    fn setup_from_underlying_sink(
901        &self,
902        cx: &mut JSContext,
903        global: &GlobalScope,
904        stream: &WritableStream,
905        underlying_sink_obj: SafeHandleObject,
906        underlying_sink: &UnderlyingSink,
907        strategy_hwm: f64,
908        strategy_size: Rc<QueuingStrategySize>,
909    ) -> Result<(), Error> {
910        // Let controller be a new WritableStreamDefaultController.
911
912        // Let startAlgorithm be an algorithm that returns undefined.
913
914        // Let writeAlgorithm be an algorithm that returns a promise resolved with undefined.
915
916        // Let closeAlgorithm be an algorithm that returns a promise resolved with undefined.
917
918        // Let abortAlgorithm be an algorithm that returns a promise resolved with undefined.
919
920        // If underlyingSinkDict["start"] exists, then set startAlgorithm to an algorithm which
921        // returns the result of invoking underlyingSinkDict["start"] with argument
922        // list « controller », exception behavior "rethrow", and callback this value underlyingSink.
923
924        // If underlyingSinkDict["write"] exists, then set writeAlgorithm to an algorithm which
925        // takes an argument chunk and returns the result of invoking underlyingSinkDict["write"]
926        // with argument list « chunk, controller » and callback this value underlyingSink.
927
928        // If underlyingSinkDict["close"] exists, then set closeAlgorithm to an algorithm which
929        // returns the result of invoking underlyingSinkDict["close"] with argument
930        // list «» and callback this value underlyingSink.
931
932        // If underlyingSinkDict["abort"] exists, then set abortAlgorithm to an algorithm which
933        // takes an argument reason and returns the result of invoking underlyingSinkDict["abort"]
934        // with argument list « reason » and callback this value underlyingSink.
935        let controller = WritableStreamDefaultController::new(
936            cx,
937            global,
938            UnderlyingSinkType::new_js(
939                underlying_sink.abort.clone(),
940                underlying_sink.start.clone(),
941                underlying_sink.close.clone(),
942                underlying_sink.write.clone(),
943            ),
944            strategy_hwm,
945            strategy_size,
946        );
947
948        // Note: this must be done before `setup`,
949        // otherwise `thisOb` is null in the start callback.
950        controller.set_underlying_sink_this_object(underlying_sink_obj);
951
952        // Perform ? SetUpWritableStreamDefaultController
953        controller.setup(cx, global, stream)
954    }
955}
956
957/// <https://streams.spec.whatwg.org/#create-writable-stream>
958#[cfg_attr(crown, expect(crown::unrooted_must_root))]
959pub(crate) fn create_writable_stream(
960    cx: &mut JSContext,
961    global: &GlobalScope,
962    writable_high_water_mark: f64,
963    writable_size_algorithm: Rc<QueuingStrategySize>,
964    underlying_sink_type: UnderlyingSinkType,
965) -> Fallible<DomRoot<WritableStream>> {
966    // Assert: ! IsNonNegativeNumber(highWaterMark) is true.
967    assert!(writable_high_water_mark >= 0.0);
968
969    // Let stream be a new WritableStream.
970    // Perform ! InitializeWritableStream(stream).
971    let stream = WritableStream::new_with_proto(cx, global, None);
972
973    // Let controller be a new WritableStreamDefaultController.
974    let controller = WritableStreamDefaultController::new(
975        cx,
976        global,
977        underlying_sink_type,
978        writable_high_water_mark,
979        writable_size_algorithm,
980    );
981
982    // Perform ? SetUpWritableStreamDefaultController(stream, controller, startAlgorithm, writeAlgorithm,
983    // closeAlgorithm, abortAlgorithm, highWaterMark, sizeAlgorithm).
984    controller.setup(cx, global, &stream)?;
985
986    // Return stream.
987    Ok(stream)
988}
989
990impl WritableStreamMethods<crate::DomTypeHolder> for WritableStream {
991    /// <https://streams.spec.whatwg.org/#ws-constructor>
992    fn Constructor(
993        cx: &mut JSContext,
994        global: &GlobalScope,
995        proto: Option<SafeHandleObject>,
996        underlying_sink: Option<*mut JSObject>,
997        strategy: &QueuingStrategy,
998    ) -> Fallible<DomRoot<WritableStream>> {
999        // If underlyingSink is missing, set it to null.
1000        rooted!(&in(cx) let underlying_sink_obj = underlying_sink.unwrap_or(ptr::null_mut()));
1001
1002        // Let underlyingSinkDict be underlyingSink,
1003        // converted to an IDL value of type UnderlyingSink.
1004        let underlying_sink_dict = if !underlying_sink_obj.is_null() {
1005            rooted!(&in(cx) let obj_val = ObjectValue(underlying_sink_obj.get()));
1006            match UnderlyingSink::new(cx, obj_val.handle()) {
1007                Ok(ConversionResult::Success(val)) => val,
1008                Ok(ConversionResult::Failure(error)) => {
1009                    return Err(Error::Type(error.into_owned()));
1010                },
1011                _ => {
1012                    return Err(Error::JSFailed);
1013                },
1014            }
1015        } else {
1016            UnderlyingSink::empty()
1017        };
1018
1019        if !underlying_sink_dict.type_.handle().is_undefined() {
1020            // If underlyingSinkDict["type"] exists, throw a RangeError exception.
1021            return Err(Error::Range(c"type is set".to_owned()));
1022        }
1023
1024        // Perform ! InitializeWritableStream(this).
1025        let stream = WritableStream::new_with_proto(cx, global, proto);
1026
1027        // Let sizeAlgorithm be ! ExtractSizeAlgorithm(strategy).
1028        let size_algorithm = extract_size_algorithm(cx, strategy);
1029
1030        // Let highWaterMark be ? ExtractHighWaterMark(strategy, 1).
1031        let high_water_mark = extract_high_water_mark(strategy, 1.0)?;
1032
1033        // Perform ? SetUpWritableStreamDefaultControllerFromUnderlyingSink(this, underlyingSink,
1034        // underlyingSinkDict, highWaterMark, sizeAlgorithm).
1035        stream.setup_from_underlying_sink(
1036            cx,
1037            global,
1038            &stream,
1039            underlying_sink_obj.handle(),
1040            &underlying_sink_dict,
1041            high_water_mark,
1042            size_algorithm,
1043        )?;
1044
1045        Ok(stream)
1046    }
1047
1048    /// <https://streams.spec.whatwg.org/#ws-locked>
1049    fn Locked(&self) -> bool {
1050        // Return ! IsWritableStreamLocked(this).
1051        self.is_locked()
1052    }
1053
1054    /// <https://streams.spec.whatwg.org/#ws-abort>
1055    fn Abort(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) -> Rc<Promise> {
1056        let global = GlobalScope::from_current_realm(cx);
1057
1058        // If ! IsWritableStreamLocked(this) is true,
1059        if self.is_locked() {
1060            // return a promise rejected with a TypeError exception.
1061            let promise = Promise::new(cx, &global);
1062            promise.reject_error(cx, Error::Type(c"Stream is locked.".to_owned()));
1063            return promise;
1064        }
1065
1066        // Return ! WritableStreamAbort(this, reason).
1067        self.abort(cx, &global, reason).into()
1068    }
1069
1070    /// <https://streams.spec.whatwg.org/#ws-close>
1071    fn Close(&self, cx: &mut CurrentRealm) -> Rc<Promise> {
1072        let global = GlobalScope::from_current_realm(cx);
1073
1074        // If ! IsWritableStreamLocked(this) is true,
1075        if self.is_locked() {
1076            // return a promise rejected with a TypeError exception.
1077            let promise = Promise::new(cx, &global);
1078            promise.reject_error(cx, Error::Type(c"Stream is locked.".to_owned()));
1079            return promise;
1080        }
1081
1082        // If ! WritableStreamCloseQueuedOrInFlight(this) is true
1083        if self.close_queued_or_in_flight() {
1084            // return a promise rejected with a TypeError exception.
1085            let promise = Promise::new(cx, &global);
1086            promise.reject_error(
1087                cx,
1088                Error::Type(c"Stream has closed queued or in-flight".to_owned()),
1089            );
1090            return promise;
1091        }
1092
1093        // Return ! WritableStreamClose(this).
1094        self.close(cx, &global).into()
1095    }
1096
1097    /// <https://streams.spec.whatwg.org/#ws-get-writer>
1098    fn GetWriter(
1099        &self,
1100        realm: &mut CurrentRealm,
1101    ) -> Result<DomRoot<WritableStreamDefaultWriter>, Error> {
1102        let global = GlobalScope::from_current_realm(realm);
1103
1104        // Return ? AcquireWritableStreamDefaultWriter(this).
1105        self.aquire_default_writer(realm, &global)
1106    }
1107}
1108
1109impl js::gc::Rootable for CrossRealmTransformWritable {}
1110
1111/// <https://streams.spec.whatwg.org/#abstract-opdef-setupcrossrealmtransformwritable>
1112/// A wrapper to handle `message` and `messageerror` events
1113/// for the port used by the transfered stream.
1114#[derive(Clone, JSTraceable, MallocSizeOf)]
1115#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
1116pub(crate) struct CrossRealmTransformWritable {
1117    /// The controller used in the algorithm.
1118    controller: Dom<WritableStreamDefaultController>,
1119
1120    /// The `backpressurePromise` used in the algorithm.
1121    #[ignore_malloc_size_of = "nested Rc"]
1122    backpressure_promise: Rc<RefCell<Option<Rc<Promise>>>>,
1123}
1124
1125impl CrossRealmTransformWritable {
1126    /// <https://streams.spec.whatwg.org/#abstract-opdef-setupcrossrealmtransformwritable>
1127    /// Add a handler for port’s message event with the following steps:
1128    pub(crate) fn handle_message(
1129        &self,
1130        cx: &mut CurrentRealm,
1131        global: &GlobalScope,
1132        message: SafeHandleValue,
1133    ) {
1134        rooted!(&in(cx) let mut value = UndefinedValue());
1135        let type_string = get_type_and_value_from_message(cx, message, value.handle_mut());
1136
1137        // If type is "pull",
1138        // Done below as the steps are the same for both types.
1139
1140        // Otherwise, if type is "error",
1141        if type_string == "error" {
1142            // Perform ! WritableStreamDefaultControllerErrorIfNeeded(controller, value).
1143            self.controller.error_if_needed(cx, value.handle(), global);
1144        }
1145
1146        let backpressure_promise = self.backpressure_promise.borrow_mut().take();
1147
1148        // Note: the below steps are for both "pull" and "error" types.
1149        // If backpressurePromise is not undefined,
1150        if let Some(promise) = backpressure_promise {
1151            // Resolve backpressurePromise with undefined.
1152            promise.resolve_native(cx, &());
1153
1154            // Set backpressurePromise to undefined.
1155            // Done above with `take`.
1156        }
1157    }
1158
1159    /// <https://streams.spec.whatwg.org/#abstract-opdef-setupcrossrealmtransformwritable>
1160    /// Add a handler for port’s messageerror event with the following steps:
1161    pub(crate) fn handle_error(
1162        &self,
1163        cx: &mut CurrentRealm,
1164        global: &GlobalScope,
1165        port: &MessagePort,
1166    ) {
1167        // Let error be a new "DataCloneError" DOMException.
1168        let error = DOMException::new(cx, global, DOMErrorName::DataCloneError);
1169        rooted!(&in(cx) let mut rooted_error = UndefinedValue());
1170        error.to_jsval(cx, rooted_error.handle_mut());
1171
1172        // Perform ! CrossRealmTransformSendError(port, error).
1173        port.cross_realm_transform_send_error(cx, rooted_error.handle());
1174
1175        // Perform ! WritableStreamDefaultControllerErrorIfNeeded(controller, error).
1176        self.controller
1177            .error_if_needed(cx, rooted_error.handle(), global);
1178
1179        // Disentangle port.
1180        global.disentangle_port(cx, port);
1181    }
1182}
1183
1184/// <https://streams.spec.whatwg.org/#ws-transfer>
1185impl Transferable for WritableStream {
1186    type Index = MessagePortIndex;
1187    type Data = MessagePortImpl;
1188
1189    /// <https://streams.spec.whatwg.org/#ref-for-transfer-steps①>
1190    fn transfer(&self, cx: &mut JSContext) -> Fallible<(MessagePortId, MessagePortImpl)> {
1191        // Step 1. If ! IsWritableStreamLocked(value) is true, throw a
1192        // "DataCloneError" DOMException.
1193        if self.is_locked() {
1194            return Err(Error::DataClone(None));
1195        }
1196
1197        let global = self.global();
1198        let mut realm = enter_auto_realm(cx, &*global);
1199        let mut realm = realm.current_realm();
1200        let cx = &mut realm;
1201
1202        // Step 2. Let port1 be a new MessagePort in the current Realm.
1203        let port_1 = MessagePort::new(cx, &global);
1204        global.track_message_port(&port_1, None);
1205
1206        // Step 3. Let port2 be a new MessagePort in the current Realm.
1207        let port_2 = MessagePort::new(cx, &global);
1208        global.track_message_port(&port_2, None);
1209
1210        // Step 4. Entangle port1 and port2.
1211        global.entangle_ports(*port_1.message_port_id(), *port_2.message_port_id());
1212
1213        // Step 5. Let readable be a new ReadableStream in the current Realm.
1214        let readable = ReadableStream::new_with_proto(cx, &global, None);
1215
1216        // Step 6. Perform ! SetUpCrossRealmTransformReadable(readable, port1).
1217        readable.setup_cross_realm_transform_readable(cx, &port_1);
1218
1219        // Step 7. Let promise be ! ReadableStreamPipeTo(readable, value, false, false, false).
1220        let promise = readable.pipe_to(cx, &global, self, false, false, false, None);
1221
1222        // Step 8. Set promise.[[PromiseIsHandled]] to true.
1223        promise.set_promise_is_handled(cx);
1224
1225        // Step 9. Set dataHolder.[[port]] to ! StructuredSerializeWithTransfer(port2, « port2 »).
1226        port_2.transfer(cx)
1227    }
1228
1229    /// <https://streams.spec.whatwg.org/#ref-for-transfer-receiving-steps①>
1230    fn transfer_receive(
1231        cx: &mut JSContext,
1232        owner: &GlobalScope,
1233        id: MessagePortId,
1234        port_impl: MessagePortImpl,
1235    ) -> Result<DomRoot<Self>, ()> {
1236        // Their transfer-receiving steps, given dataHolder and value, are:
1237        // Note: dataHolder is used in `structuredclone.rs`, and value is created here.
1238        let value = WritableStream::new_with_proto(cx, owner, None);
1239
1240        // Step 1. Let deserializedRecord be !
1241        // StructuredDeserializeWithTransfer(dataHolder.[[port]], the current
1242        // Realm).
1243        // Done with the `Deserialize` derive of `MessagePortImpl`.
1244
1245        // Step 2. Let port be deserializedRecord.[[Deserialized]].
1246        let transferred_port = MessagePort::transfer_receive(cx, owner, id, port_impl)?;
1247
1248        // Step 3. Perform ! SetUpCrossRealmTransformWritable(value, port).
1249        value.setup_cross_realm_transform_writable(cx, &transferred_port);
1250        Ok(value)
1251    }
1252
1253    /// Note: we are relying on the port transfer, so the data returned here are related to the port.
1254    fn serialized_storage<'a>(
1255        data: StructuredData<'a, '_>,
1256    ) -> &'a mut Option<FxHashMap<MessagePortId, Self::Data>> {
1257        match data {
1258            StructuredData::Reader(r) => &mut r.port_impls,
1259            StructuredData::Writer(w) => &mut w.ports,
1260        }
1261    }
1262}