Skip to main content

ferrijs_std/stream_web/readable/stream/
mod.rs

1use std::{cell::OnceCell, panic, rc::Rc};
2
3use crate::stream_web::{
4    queuing_strategy::{QueuingStrategy, SizeAlgorithm},
5    readable::{
6        byob_reader::{ReadableStreamBYOBReader, ReadableStreamReadIntoRequest, ViewBytes},
7        byte_controller::{ReadableByteStreamController, ReadableByteStreamControllerClass},
8        controller::{ReadableStreamController, ReadableStreamControllerClass},
9        default_controller::{
10            ReadableStreamDefaultController, ReadableStreamDefaultControllerOwned,
11        },
12        default_reader::{ReadableStreamDefaultReader, ReadableStreamReadRequest},
13        iterator::{IteratorKind, IteratorRecord, ReadableStreamAsyncIterator},
14        objects::{
15            ReadableStreamBYOBObjects, ReadableStreamClassObjects,
16            ReadableStreamDefaultReaderObjects, ReadableStreamObjects,
17        },
18        reader::{
19            ReadableStreamReader, ReadableStreamReaderClass, ReadableStreamReaderOwned,
20            UndefinedReader,
21        },
22    },
23    readable_writable_pair::ReadableWritablePair,
24    utils::{
25        promise::{
26            promise_rejected_catch, promise_rejected_with, promise_rejected_with_constructor,
27            promise_resolved_with, upon_promise_fulfilment, with_promise_result,
28            PromisePrimordials,
29        },
30        UnwrapOrUndefined, ValueOrUndefined,
31    },
32    writable::WritableStreamOwned,
33};
34
35use pipe::StreamPipeOptions;
36use source::UnderlyingSource;
37
38pub use algorithms::{CancelAlgorithm, PullAlgorithm, StartAlgorithm};
39use crate::utils::{
40    option::{Null, NullableOpt, Undefined},
41    primordials::{BasePrimordials, Primordial},
42    result::ResultExt,
43};
44use rquickjs::{
45    atom::PredefinedAtom,
46    class::{OwnedBorrowMut, Trace},
47    function::Constructor,
48    prelude::{List, Opt, This},
49    Class, Coerced, Ctx, Error, Exception, FromJs, Function, IntoJs, JsLifetime, Object, Promise,
50    Result, Value,
51};
52
53pub mod algorithms;
54mod pipe;
55pub(super) mod source;
56mod tee;
57
58/// Acquire a default reader for the stream, locking it. Subsequent
59/// `getReader()` calls from JS will throw per spec.
60pub fn lock_readable_stream<'js>(
61    ctx: Ctx<'js>,
62    stream: Class<'js, ReadableStream<'js>>,
63) -> Result<()> {
64    let owned = rquickjs::class::OwnedBorrowMut::from_class(stream);
65    super::reader::ReadableStreamReaderClass::acquire_readable_stream_default_reader(ctx, owned)?;
66    Ok(())
67}
68
69/// Fast-path drain for a default-controller ReadableStream whose queue holds
70/// all the data synchronously (e.g. the stream was enqueued in `start()` and
71/// then closed). Bypasses the JS reader + Promise machinery, so user code
72/// that poisons `Object.prototype.then` cannot swap the streamed chunks
73/// (WPT `response-stream-with-broken-then`).
74///
75/// Returns `Some(chunks)` if the fast path applied, `None` otherwise (stream
76/// locked, disturbed, has pending pull, byte controller, not yet closed,
77/// etc). Sets `disturbed = true` on success.
78pub fn try_sync_drain_closed_stream<'js>(
79    stream: &Class<'js, ReadableStream<'js>>,
80) -> Option<Vec<rquickjs::Value<'js>>> {
81    use super::controller::ReadableStreamControllerClass;
82    use super::default_controller::ReadableStreamDefaultController;
83    use super::stream::ReadableStreamState;
84    use rquickjs::class::OwnedBorrowMut;
85
86    let mut stream_ref = stream.try_borrow_mut().ok()?;
87    if stream_ref.disturbed || stream_ref.is_readable_stream_locked() {
88        return None;
89    }
90    // Stream state must be Readable (not Errored). Closed would also be OK
91    // but then the queue should already be empty.
92    if !matches!(stream_ref.state, ReadableStreamState::Readable) {
93        return None;
94    }
95    let controller_class = match &stream_ref.controller {
96        ReadableStreamControllerClass::ReadableStreamDefaultController(c) => c.clone(),
97        _ => return None,
98    };
99    let mut controller: OwnedBorrowMut<'js, ReadableStreamDefaultController<'js>> =
100        OwnedBorrowMut::try_from_class(controller_class).ok()?;
101    // Only fast-path when close has been requested — otherwise there could
102    // be more data coming via `pull()` that we'd miss.
103    if !controller.close_requested {
104        return None;
105    }
106    let mut chunks = Vec::with_capacity(controller.container.queue.len());
107    while !controller.container.queue.is_empty() {
108        chunks.push(controller.container.dequeue_value());
109    }
110    stream_ref.disturbed = true;
111    // Transition the stream to Closed now that its queue is drained, so that
112    // later consumers see a consistent state.
113    stream_ref.state = ReadableStreamState::Closed;
114    Some(chunks)
115}
116
117/// Tee a ReadableStream into two branches. The stream must not be locked or disturbed.
118pub fn tee_readable_stream<'js>(
119    ctx: Ctx<'js>,
120    stream: Class<'js, ReadableStream<'js>>,
121) -> Result<(
122    Class<'js, ReadableStream<'js>>,
123    Class<'js, ReadableStream<'js>>,
124)> {
125    {
126        let stream_ref = stream.borrow();
127        if stream_ref.disturbed {
128            return Err(Exception::throw_type(
129                &ctx,
130                "Cannot tee a disturbed ReadableStream",
131            ));
132        }
133        if stream_ref.is_readable_stream_locked() {
134            return Err(Exception::throw_type(
135                &ctx,
136                "Cannot tee a locked ReadableStream",
137            ));
138        }
139    }
140    let owned = OwnedBorrowMut::from_class(stream);
141    let objects = ReadableStreamObjects::from_stream(owned);
142    ReadableStream::readable_stream_tee(ctx, objects)
143}
144
145#[rquickjs::class]
146#[derive(JsLifetime)]
147pub struct ReadableStream<'js> {
148    pub controller: ReadableStreamControllerClass<'js>,
149    pub disturbed: bool,
150    pub state: ReadableStreamState<'js>,
151    pub(crate) reader: Option<ReadableStreamReaderClass<'js>>,
152    pub(crate) promise_primordials: PromisePrimordials<'js>,
153    pub(crate) constructor_type_error: Constructor<'js>,
154    pub(crate) constructor_range_error: Constructor<'js>,
155    pub(crate) function_array_buffer_is_view: Function<'js>,
156}
157
158impl<'js> Trace<'js> for ReadableStream<'js> {
159    fn trace<'a>(&self, tracer: rquickjs::class::Tracer<'a, 'js>) {
160        self.controller.trace(tracer);
161        self.state.trace(tracer);
162        self.reader.trace(tracer);
163
164        self.promise_primordials.trace(tracer);
165        self.constructor_type_error.trace(tracer);
166        self.constructor_range_error.trace(tracer);
167        self.function_array_buffer_is_view.trace(tracer);
168    }
169}
170
171pub(crate) type ReadableStreamClass<'js> = Class<'js, ReadableStream<'js>>;
172pub(crate) type ReadableStreamOwned<'js> = OwnedBorrowMut<'js, ReadableStream<'js>>;
173
174#[derive(Debug, Trace, Clone, JsLifetime)]
175pub enum ReadableStreamState<'js> {
176    Readable,
177    Closed,
178    Errored(Value<'js>),
179}
180
181#[rquickjs::methods(rename_all = "camelCase")]
182impl<'js> ReadableStream<'js> {
183    // Streams Spec: 4.2.4: https://streams.spec.whatwg.org/#rs-prototype
184    // constructor(optional object underlyingSource, optional QueuingStrategy strategy = {});
185    #[qjs(constructor)]
186    fn new(
187        ctx: Ctx<'js>,
188        underlying_source: Opt<Undefined<Object<'js>>>,
189        queuing_strategy: Opt<Undefined<QueuingStrategy<'js>>>,
190    ) -> Result<Class<'js, Self>> {
191        // If underlyingSource is missing, set it to null.
192        let underlying_source = Null(underlying_source.0);
193
194        // Let underlyingSourceDict be underlyingSource, converted to an IDL value of type UnderlyingSource.
195        let underlying_source_dict = match underlying_source {
196            Null(None) | Null(Some(Undefined(None))) => UnderlyingSource::default(),
197            Null(Some(Undefined(Some(ref obj)))) => UnderlyingSource::from_object(obj.clone())?,
198        };
199
200        let promise_primordials = PromisePrimordials::get(&ctx)?.clone();
201        let base_primordials = BasePrimordials::get(&ctx)?;
202
203        let stream_class = Class::instance(
204            ctx.clone(),
205            Self {
206                // Set stream.[[state]] to "readable".
207                state: ReadableStreamState::Readable,
208                // Set stream.[[reader]] and stream.[[storedError]] to undefined.
209                reader: None,
210                // Set stream.[[disturbed]] to false.
211                disturbed: false,
212                controller: ReadableStreamControllerClass::Uninitialised,
213                constructor_type_error: base_primordials.constructor_type_error.clone(),
214                constructor_range_error: base_primordials.constructor_range_error.clone(),
215                function_array_buffer_is_view: base_primordials
216                    .function_array_buffer_is_view
217                    .clone(),
218                promise_primordials,
219            },
220        )?;
221        drop(base_primordials);
222        let stream = OwnedBorrowMut::from_class(stream_class.clone());
223        let queuing_strategy = queuing_strategy.0.and_then(|qs| qs.0);
224
225        match underlying_source_dict.r#type {
226            // If underlyingSourceDict["type"] is "bytes":
227            Some(ReadableStreamType::Bytes) => {
228                // If strategy["size"] exists, throw a RangeError exception.
229                if queuing_strategy
230                    .as_ref()
231                    .and_then(|qs| qs.size.as_ref())
232                    .is_some()
233                {
234                    return Err(Exception::throw_range(
235                        &ctx,
236                        "The strategy for a byte stream cannot have a size function",
237                    ));
238                }
239                // Let highWaterMark be ? ExtractHighWaterMark(strategy, 0).
240                let high_water_mark =
241                    QueuingStrategy::extract_high_water_mark(&ctx, queuing_strategy, 0.0)?;
242
243                // Perform ? SetUpReadableByteStreamControllerFromUnderlyingSource(this, underlyingSource, underlyingSourceDict, highWaterMark).
244                ReadableByteStreamController::set_up_readable_byte_stream_controller_from_underlying_source(
245                    &ctx,
246                    stream,
247                    underlying_source,
248                    underlying_source_dict,
249                    high_water_mark,
250                )?;
251            },
252            // Otherwise (no type, or "owning" which we treat as a default
253            // controller that also accepts the `transfer` enqueue option):
254            None | Some(ReadableStreamType::Owning) => {
255                let is_owning_type = matches!(
256                    underlying_source_dict.r#type,
257                    Some(ReadableStreamType::Owning)
258                );
259                // Let sizeAlgorithm be ! ExtractSizeAlgorithm(strategy).
260                let size_algorithm =
261                    QueuingStrategy::extract_size_algorithm(queuing_strategy.as_ref());
262
263                // Let highWaterMark be ? ExtractHighWaterMark(strategy, 1).
264                let high_water_mark =
265                    QueuingStrategy::extract_high_water_mark(&ctx, queuing_strategy, 1.0)?;
266
267                // Perform ? SetUpReadableStreamDefaultControllerFromUnderlyingSource(this, underlyingSource, underlyingSourceDict, highWaterMark, sizeAlgorithm).
268                ReadableStreamDefaultController::set_up_readable_stream_default_controller_from_underlying_source(
269                    ctx,
270                    stream,
271                    underlying_source,
272                    underlying_source_dict,
273                    high_water_mark,
274                    size_algorithm,
275                    is_owning_type,
276                )?;
277            },
278        }
279
280        Ok(stream_class)
281    }
282
283    // static ReadableStream from(any asyncIterable);
284    #[qjs(static)]
285    fn from(ctx: Ctx<'js>, async_iterable: Value<'js>) -> Result<Class<'js, Self>> {
286        // Return ? ReadableStreamFromIterable(asyncIterable).
287        Self::readable_stream_from_iterable(&ctx, async_iterable)
288    }
289
290    // readonly attribute boolean locked;
291    #[qjs(get)]
292    fn locked(&self) -> bool {
293        // Return ! IsReadableStreamLocked(this).
294        self.is_readable_stream_locked()
295    }
296
297    // Internal property for checking if stream has been read from
298    #[qjs(get)]
299    fn disturbed(&self) -> bool {
300        self.disturbed
301    }
302
303    // Promise<undefined> cancel(optional any reason);
304    fn cancel(
305        ctx: Ctx<'js>,
306        stream: This<OwnedBorrowMut<'js, Self>>,
307        reason: Opt<Value<'js>>,
308    ) -> Result<Promise<'js>> {
309        // If ! IsReadableStreamLocked(this) is true, return a promise rejected with a TypeError exception.
310        if stream.is_readable_stream_locked() {
311            return promise_rejected_with_constructor(
312                &stream.constructor_type_error,
313                &stream.promise_primordials,
314                "Cannot cancel a stream that already has a reader",
315            );
316        }
317
318        let objects = ReadableStreamObjects::from_stream(stream.0).refresh_reader();
319
320        let (promise, _) =
321            Self::readable_stream_cancel(ctx.clone(), objects, reason.0.unwrap_or_undefined(&ctx))?;
322        Ok(promise)
323    }
324
325    // ReadableStreamReader getReader(optional ReadableStreamGetReaderOptions options = {});
326    fn get_reader(
327        ctx: Ctx<'js>,
328        stream: This<OwnedBorrowMut<'js, Self>>,
329        options: Opt<Option<ReadableStreamGetReaderOptions>>,
330    ) -> Result<ReadableStreamReaderClass<'js>> {
331        // If options["mode"] does not exist, return ? AcquireReadableStreamDefaultReader(this).
332        let reader = match options.0 {
333            None | Some(None | Some(ReadableStreamGetReaderOptions { mode: None })) => {
334                let (_, reader) =
335                    ReadableStreamReaderClass::acquire_readable_stream_default_reader(
336                        ctx.clone(),
337                        stream.0,
338                    )?;
339                reader.into()
340            },
341            // Return ? AcquireReadableStreamBYOBReader(this).
342            Some(Some(ReadableStreamGetReaderOptions {
343                mode: Some(ReadableStreamReaderMode::Byob),
344            })) => {
345                let (_, reader) = ReadableStreamReaderClass::acquire_readable_stream_byob_reader(
346                    ctx.clone(),
347                    stream.0,
348                )?;
349                reader.into()
350            },
351        };
352
353        Ok(reader)
354    }
355
356    // ReadableStream pipeThrough(ReadableWritablePair transform, optional StreamPipeOptions options = {});
357    fn pipe_through(
358        ctx: Ctx<'js>,
359        stream: This<OwnedBorrowMut<'js, Self>>,
360        transform: ReadableWritablePair<'js>,
361        options: NullableOpt<StreamPipeOptions<'js>>,
362    ) -> Result<ReadableStreamClass<'js>> {
363        // If ! IsReadableStreamLocked(this) is true, throw a TypeError exception.
364        if stream.is_readable_stream_locked() {
365            return Err(Exception::throw_type(
366                &ctx,
367                "ReadableStream.prototype.pipeThrough cannot be used on a locked ReadableStream",
368            ));
369        }
370
371        let readable_class = transform.readable.clone();
372        let writable = OwnedBorrowMut::from_class(transform.writable);
373
374        // If ! IsWritableStreamLocked(transform["writable"]) is true, throw a TypeError exception.
375        if writable.is_writable_stream_locked() {
376            return Err(Exception::throw_type(
377                &ctx,
378                "ReadableStream.prototype.pipeThrough cannot be used on a locked WritableStream",
379            ));
380        }
381
382        // Let signal be options["signal"] if it exists, or undefined otherwise.
383        let options = options.0.unwrap_or_default();
384
385        // Let promise be ! ReadableStreamPipeTo(this, transform["writable"], options["preventClose"], options["preventAbort"], options["preventCancel"], signal).
386        let promise = ReadableStream::readable_stream_pipe_to(
387            ctx.clone(),
388            stream.0,
389            writable,
390            options.prevent_close,
391            options.prevent_abort,
392            options.prevent_cancel,
393            options.signal,
394        )?;
395
396        // Set promise.[[PromiseIsHandled]] to true.
397        let () = promise
398            .catch()?
399            .call((This(promise.clone()), Function::new(ctx, || {})))?;
400
401        // Return transform["readable"].
402        Ok(readable_class)
403    }
404
405    // Promise<undefined> pipeTo(WritableStream destination, optional StreamPipeOptions options = {});
406    fn pipe_to(
407        ctx: Ctx<'js>,
408        stream: This<Value<'js>>,
409        destination: Value<'js>,
410        options: NullableOpt<Value<'js>>,
411    ) -> Result<Promise<'js>> {
412        with_promise_result(&ctx, || {
413            let stream =
414                ReadableStreamOwned::from_class(Class::from_value(&stream.0).or_throw_type(
415                    &ctx,
416                    "'pipeTo' called on an object that is not a valid instance of ReadableStream.",
417                )?);
418
419            let options = match options.0 {
420                Some(options) => Some(StreamPipeOptions::from_js(&ctx, options)?),
421                None => None,
422            };
423
424            // If ! IsReadableStreamLocked(this) is true, return a promise rejected with a TypeError exception.
425            if stream.is_readable_stream_locked() {
426                return promise_rejected_with_constructor(
427                    &stream.constructor_type_error,
428                    &stream.promise_primordials,
429                    "ReadableStream.prototype.pipeTo cannot be used on a locked ReadableStream",
430                );
431            }
432
433            let destination = WritableStreamOwned::from_class(
434                Class::from_value(&destination).or_throw_type(&ctx,"'pipeTo' instructed to pipe to an object that is not a valid instance of WritableStream.")?,
435            );
436
437            // If ! IsWritableStreamLocked(destination) is true, return a promise rejected with a TypeError exception.
438            if destination.is_writable_stream_locked() {
439                return promise_rejected_with_constructor(
440                    &stream.constructor_type_error,
441                    &stream.promise_primordials,
442                    "ReadableStream.prototype.pipeTo cannot be used on a locked WritableStream",
443                );
444            }
445
446            // Let signal be options["signal"] if it exists, or undefined otherwise.
447            let options = options.unwrap_or_default();
448
449            // Return ! ReadableStreamPipeTo(this, destination, options["preventClose"], options["preventAbort"], options["preventCancel"], signal).
450            Self::readable_stream_pipe_to(
451                ctx.clone(),
452                stream,
453                destination,
454                options.prevent_close,
455                options.prevent_abort,
456                options.prevent_cancel,
457                options.signal,
458            )
459        })
460    }
461
462    // sequence<ReadableStream> tee();
463    fn tee(
464        ctx: Ctx<'js>,
465        stream: This<OwnedBorrowMut<'js, Self>>,
466    ) -> Result<List<(Class<'js, Self>, Class<'js, Self>)>> {
467        Ok(List(Self::readable_stream_tee(
468            ctx,
469            ReadableStreamObjects::from_stream(stream.0),
470        )?))
471    }
472
473    #[qjs(rename = PredefinedAtom::SymbolAsyncIterator)]
474    fn async_iterate(
475        ctx: Ctx<'js>,
476        stream: This<OwnedBorrowMut<'js, Self>>,
477    ) -> Result<Class<'js, ReadableStreamAsyncIterator<'js>>> {
478        Self::values(ctx, stream, Opt(None))
479    }
480
481    fn values(
482        ctx: Ctx<'js>,
483        stream: This<OwnedBorrowMut<'js, Self>>,
484        arg: Opt<Object<'js>>,
485    ) -> Result<Class<'js, ReadableStreamAsyncIterator<'js>>> {
486        // Let reader be ? AcquireReadableStreamDefaultReader(stream).
487        let (stream, reader) = ReadableStreamReaderClass::acquire_readable_stream_default_reader(
488            ctx.clone(),
489            stream.0,
490        )?;
491
492        // Let preventCancel be args[0]["preventCancel"].
493        let prevent_cancel = match arg.0 {
494            None => false,
495            Some(arg) => matches!(arg.get_value_or_undefined("preventCancel")?, Some(true)),
496        };
497
498        let promise_primordials = stream.promise_primordials.clone();
499        let controller = stream.controller.clone();
500
501        ReadableStreamAsyncIterator::new(
502            ctx,
503            ReadableStreamClassObjects {
504                stream: stream.into_inner(),
505                controller,
506                reader,
507            },
508            promise_primordials,
509            prevent_cancel,
510        )
511    }
512}
513
514impl<'js> ReadableStream<'js> {
515    pub(super) fn readable_stream_error<
516        C: ReadableStreamController<'js>,
517        R: ReadableStreamReader<'js>,
518    >(
519        // Let reader be stream.[[reader]].
520        mut objects: ReadableStreamObjects<'js, C, R>,
521        e: Value<'js>,
522    ) -> Result<ReadableStreamObjects<'js, C, R>> {
523        // Set stream.[[state]] to "errored".
524        // Set stream.[[storedError]] to e.
525        objects.stream.state = ReadableStreamState::Errored(e.clone());
526
527        objects = objects.with_reader(
528            // If reader implements ReadableStreamDefaultReader,
529            |mut objects| {
530                // Reject reader.[[closedPromise]] with e.
531                objects.reader
532                    .generic
533                    .closed_promise
534                    .reject(e.clone())?;
535
536                // Set reader.[[closedPromise]].[[PromiseIsHandled]] to true.
537                objects.reader.generic.closed_promise.set_is_handled()?;
538
539                // Perform ! ReadableStreamDefaultReaderErrorReadRequests(reader, e).
540                objects = ReadableStreamDefaultReader::readable_stream_default_reader_error_read_requests(
541                        objects, e.clone(),
542                )?;
543                Ok(objects)
544        },
545            // Otherwise,
546            |mut objects| {
547                // Reject reader.[[closedPromise]] with e.
548                objects.reader
549                    .generic
550                    .closed_promise
551                    .reject(e.clone())?;
552
553                // Set reader.[[closedPromise]].[[PromiseIsHandled]] to true.
554                objects.reader.generic.closed_promise.set_is_handled()?;
555
556                // Perform ! ReadableStreamBYOBReaderErrorReadIntoRequests(reader, e).
557                objects = ReadableStreamBYOBReader::readable_stream_byob_reader_error_read_into_requests(
558                    objects, e.clone(),
559                )?;
560
561                Ok(objects)
562            },
563        // If reader is undefined, return.
564        Ok)?;
565
566        Ok(objects)
567    }
568
569    pub(super) fn readable_stream_get_num_read_requests(
570        reader: &ReadableStreamDefaultReader,
571    ) -> usize {
572        reader.read_requests.len()
573    }
574
575    pub(super) fn readable_stream_get_num_read_into_requests(
576        reader: &ReadableStreamBYOBReader,
577    ) -> usize {
578        reader.read_into_requests.len()
579    }
580
581    pub(super) fn readable_stream_fulfill_read_request<C: ReadableStreamController<'js>>(
582        ctx: &Ctx<'js>,
583        // Let reader be stream.[[reader]].
584        mut objects: ReadableStreamDefaultReaderObjects<'js, C>,
585        chunk: Value<'js>,
586        done: bool,
587    ) -> Result<ReadableStreamDefaultReaderObjects<'js, C>> {
588        // Let readRequest be reader.[[readRequests]][0].
589        // Remove readRequest from reader.[[readRequests]].
590        let read_request = objects
591            .reader
592            .read_requests
593            .pop_front()
594            .expect("ReadableStreamFulfillReadRequest called with empty readRequests");
595
596        if done {
597            // If done is true, perform readRequest’s close steps.
598            read_request.close_steps_typed(ctx, objects)
599        } else {
600            // Otherwise, perform readRequest’s chunk steps, given chunk.
601            read_request.chunk_steps_typed(objects, chunk)
602        }
603    }
604
605    pub(super) fn readable_stream_fulfill_read_into_request(
606        ctx: &Ctx<'js>,
607        mut objects: ReadableStreamBYOBObjects<'js>,
608        chunk: ViewBytes<'js>,
609        done: bool,
610    ) -> Result<ReadableStreamBYOBObjects<'js>> {
611        // Let readIntoRequest be reader.[[readIntoRequests]][0].
612        // Remove readIntoRequest from reader.[[readIntoRequests]].
613        let read_into_request = objects
614            .reader
615            .read_into_requests
616            .pop_front()
617            .expect("ReadableStreamFulfillReadIntoRequest called with empty readIntoRequests");
618
619        if done {
620            // If done is true, perform readIntoRequest’s close steps, given chunk.
621            read_into_request.close_steps(objects, chunk.into_js(ctx)?)
622        } else {
623            // Otherwise, perform readIntoRequest’s chunk steps, given chunk.
624            read_into_request.chunk_steps(objects, chunk.into_js(ctx)?)
625        }
626    }
627
628    pub(super) fn readable_stream_close<
629        C: ReadableStreamController<'js>,
630        R: ReadableStreamReader<'js>,
631    >(
632        ctx: Ctx<'js>,
633        // Let reader be stream.[[reader]].
634        mut objects: ReadableStreamObjects<'js, C, R>,
635    ) -> Result<ReadableStreamObjects<'js, C, R>> {
636        // Set stream.[[state]] to "closed".
637        objects.stream.state = ReadableStreamState::Closed;
638
639        objects.with_reader(
640            |mut objects| {
641                // Resolve reader.[[closedPromise]] with undefined.
642                objects.reader.generic.closed_promise.resolve_undefined()?;
643
644                // If reader implements ReadableStreamDefaultReader,
645                // Let readRequests be reader.[[readRequests]].
646                // Set reader.[[readRequests]] to an empty list.
647                let read_requests = objects.reader.read_requests.split_off(0);
648
649                // For each readRequest of readRequests,
650                for read_request in read_requests {
651                    // Perform readRequest’s close steps.
652                    objects = read_request.close_steps_typed(&ctx, objects)?;
653                }
654
655                Ok(objects)
656            },
657            |objects| {
658                objects.reader.generic.closed_promise.resolve_undefined()?;
659
660                Ok(objects)
661            },
662            // If reader is undefined, return.
663            Ok,
664        )
665    }
666
667    pub fn is_readable_stream_locked(&self) -> bool {
668        // If stream.[[reader]] is undefined, return false.
669        if self.reader.is_none() {
670            return false;
671        }
672        // Return true.
673        true
674    }
675
676    pub(super) fn readable_stream_add_read_request(
677        &mut self,
678        reader: &mut ReadableStreamDefaultReader<'js>,
679        read_request: impl ReadableStreamReadRequest<'js> + 'js,
680    ) {
681        reader.read_requests.push_back(Box::new(read_request));
682    }
683
684    pub(super) fn readable_stream_cancel<
685        C: ReadableStreamController<'js>,
686        R: ReadableStreamReader<'js>,
687    >(
688        ctx: Ctx<'js>,
689        mut objects: ReadableStreamObjects<'js, C, R>,
690        reason: Value<'js>,
691    ) -> Result<(Promise<'js>, ReadableStreamObjects<'js, C, R>)> {
692        // Set stream.[[disturbed]] to true.
693        objects.stream.disturbed = true;
694
695        match objects.stream.state {
696            // If stream.[[state]] is "closed", return a promise resolved with undefined.
697            ReadableStreamState::Closed => Ok((
698                // wpt tests expect that this is a new promise every time so we can't duplicate the primordial promise_resolved_with_undefined
699                promise_resolved_with(
700                    &ctx,
701                    &objects.stream.promise_primordials,
702                    Ok(Value::new_undefined(ctx.clone())),
703                )?,
704                objects,
705            )),
706            // If stream.[[state]] is "errored", return a promise rejected with stream.[[storedError]].
707            ReadableStreamState::Errored(ref stored_error) => Ok((
708                promise_rejected_with(&objects.stream.promise_primordials, stored_error.clone())?,
709                objects,
710            )),
711            ReadableStreamState::Readable => {
712                // Perform ! ReadableStreamClose(stream).
713                objects = ReadableStream::readable_stream_close(ctx.clone(), objects)?;
714                // Let reader be stream.[[reader]].
715                // If reader is not undefined and reader implements ReadableStreamBYOBReader,
716
717                objects = objects.with_reader(
718                    Ok,
719                    |mut objects| {
720                        // Let readIntoRequests be reader.[[readIntoRequests]].
721                        // Set reader.[[readIntoRequests]] to an empty list.
722                        let read_into_requests = objects.reader.read_into_requests.split_off(0);
723                        // For each readIntoRequest of readIntoRequests,
724                        for read_into_request in read_into_requests {
725                            // Perform readIntoRequest’s close steps, given undefined.
726                            objects = read_into_request
727                                .close_steps(objects, Value::new_undefined(ctx.clone()))?;
728                        }
729
730                        Ok(objects)
731                    },
732                    Ok,
733                )?;
734
735                // Let sourceCancelPromise be ! stream.[[controller]].[[CancelSteps]](reason).
736                let (source_cancel_promise, objects) = C::cancel_steps(&ctx, objects, reason)?;
737
738                // Return the result of reacting to sourceCancelPromise with a fulfillment step that returns undefined.
739                let promise = upon_promise_fulfilment(ctx, source_cancel_promise, |_, ()| {
740                    Ok(rquickjs::Undefined)
741                })?;
742
743                Ok((promise, objects))
744            },
745        }
746    }
747
748    pub(super) fn readable_stream_add_read_into_request(
749        reader: &mut ReadableStreamBYOBReader<'js>,
750        read_request: impl ReadableStreamReadIntoRequest<'js> + 'js,
751    ) {
752        // Append readRequest to stream.[[reader]].[[readIntoRequests]].
753        reader.read_into_requests.push_back(Box::new(read_request))
754    }
755
756    // CreateReadableStream(startAlgorithm, pullAlgorithm, cancelAlgorithm[, highWaterMark, [, sizeAlgorithm]]) performs the following steps:
757    pub(crate) fn create_readable_stream(
758        ctx: Ctx<'js>,
759        start_algorithm: StartAlgorithm<'js>,
760        pull_algorithm: PullAlgorithm<'js>,
761        cancel_algorithm: CancelAlgorithm<'js>,
762        high_water_mark: Option<f64>,
763        size_algorithm: Option<SizeAlgorithm<'js>>,
764    ) -> Result<
765        ReadableStreamClassObjects<'js, ReadableStreamDefaultControllerOwned<'js>, UndefinedReader>,
766    > {
767        // If highWaterMark was not passed, set it to 1.
768        let high_water_mark = high_water_mark.unwrap_or(1.0);
769
770        // If sizeAlgorithm was not passed, set it to an algorithm that returns 1.
771        let size_algorithm = size_algorithm.unwrap_or(SizeAlgorithm::AlwaysOne);
772
773        let base_primordials = BasePrimordials::get(&ctx)?;
774
775        // Let stream be a new ReadableStream.
776        let stream_class = Class::instance(
777            ctx.clone(),
778            Self {
779                // Set stream.[[state]] to "readable".
780                state: ReadableStreamState::Readable,
781                // Set stream.[[reader]] and stream.[[storedError]] to undefined.
782                reader: None,
783                // Set stream.[[disturbed]] to false.
784                disturbed: false,
785                controller: ReadableStreamControllerClass::Uninitialised,
786                promise_primordials: PromisePrimordials::get(&ctx)?.clone(),
787                constructor_range_error: base_primordials.constructor_range_error.clone(),
788                constructor_type_error: base_primordials.constructor_type_error.clone(),
789                function_array_buffer_is_view: base_primordials
790                    .function_array_buffer_is_view
791                    .clone(),
792            },
793        )?;
794        drop(base_primordials);
795
796        // Perform ? SetUpReadableStreamDefaultController(stream, controller, startAlgorithm, pullAlgorithm, cancelAlgorithm, highWaterMark, sizeAlgorithm).
797        let controller_class =
798            ReadableStreamDefaultController::set_up_readable_stream_default_controller(
799                ctx,
800                OwnedBorrowMut::from_class(stream_class.clone()),
801                start_algorithm,
802                pull_algorithm,
803                cancel_algorithm,
804                high_water_mark,
805                size_algorithm,
806                false, // not owning-type; Rust-side streams never set it
807            )?;
808
809        // Return stream.
810        Ok(ReadableStreamClassObjects {
811            stream: stream_class,
812            controller: controller_class,
813            reader: UndefinedReader,
814        })
815    }
816
817    /// Create a ReadableStream from Rust pull/cancel algorithms
818    pub fn from_pull_algorithm(
819        ctx: Ctx<'js>,
820        pull_algorithm: PullAlgorithm<'js>,
821        cancel_algorithm: CancelAlgorithm<'js>,
822    ) -> Result<Class<'js, Self>> {
823        Self::from_pull_algorithm_with_options(ctx, pull_algorithm, cancel_algorithm, None)
824    }
825
826    /// Create a ReadableStream from Rust pull/cancel algorithms with custom highWaterMark
827    pub fn from_pull_algorithm_with_options(
828        ctx: Ctx<'js>,
829        pull_algorithm: PullAlgorithm<'js>,
830        cancel_algorithm: CancelAlgorithm<'js>,
831        high_water_mark: Option<f64>,
832    ) -> Result<Class<'js, Self>> {
833        Ok(Self::create_readable_stream(
834            ctx,
835            StartAlgorithm::ReturnUndefined,
836            pull_algorithm,
837            cancel_algorithm,
838            high_water_mark,
839            None,
840        )?
841        .stream)
842    }
843
844    /// Create a byte-source ReadableStream (i.e. `type: "bytes"`) from Rust
845    /// pull/cancel algorithms. BYOB readers can attach to the returned
846    /// stream, and the pull algorithm receives a byte controller so it can
847    /// enqueue `Uint8Array` chunks that stream directly into BYOB reads
848    /// without copying.
849    pub fn from_byte_pull_algorithm(
850        ctx: Ctx<'js>,
851        pull_algorithm: PullAlgorithm<'js>,
852        cancel_algorithm: CancelAlgorithm<'js>,
853    ) -> Result<Class<'js, Self>> {
854        let (stream, _controller) = Self::create_readable_byte_stream(
855            ctx,
856            StartAlgorithm::ReturnUndefined,
857            pull_algorithm,
858            cancel_algorithm,
859        )?;
860        Ok(stream)
861    }
862
863    // CreateReadableByteStream(startAlgorithm, pullAlgorithm, cancelAlgorithm) performs the following steps:
864    pub fn create_readable_byte_stream(
865        ctx: Ctx<'js>,
866        start_algorithm: StartAlgorithm<'js>,
867        pull_algorithm: PullAlgorithm<'js>,
868        cancel_algorithm: CancelAlgorithm<'js>,
869    ) -> Result<(Class<'js, Self>, ReadableByteStreamControllerClass<'js>)> {
870        let base_primordials = BasePrimordials::get(&ctx)?;
871
872        // Let stream be a new ReadableStream.
873        let stream_class = Class::instance(
874            ctx.clone(),
875            Self {
876                // Set stream.[[state]] to "readable".
877                state: ReadableStreamState::Readable,
878                // Set stream.[[reader]] and stream.[[storedError]] to undefined.
879                reader: None,
880                // Set stream.[[disturbed]] to false.
881                disturbed: false,
882                controller: ReadableStreamControllerClass::Uninitialised,
883                promise_primordials: PromisePrimordials::get(&ctx)?.clone(),
884                constructor_type_error: base_primordials.constructor_type_error.clone(),
885                constructor_range_error: base_primordials.constructor_range_error.clone(),
886                function_array_buffer_is_view: base_primordials
887                    .function_array_buffer_is_view
888                    .clone(),
889            },
890        )?;
891        drop(base_primordials);
892
893        // Perform ? SetUpReadableStreamDefaultController(stream, controller, startAlgorithm, pullAlgorithm, cancelAlgorithm, highWaterMark, sizeAlgorithm).
894        let controller_class =
895            ReadableByteStreamController::set_up_readable_byte_stream_controller(
896                ctx,
897                OwnedBorrowMut::from_class(stream_class.clone()),
898                start_algorithm,
899                pull_algorithm,
900                cancel_algorithm,
901                0.0,
902                None,
903            )?;
904
905        // Return stream.
906        Ok((stream_class, controller_class))
907    }
908
909    fn readable_stream_from_iterable(
910        ctx: &Ctx<'js>,
911        async_iterable: Value<'js>,
912    ) -> Result<Class<'js, Self>> {
913        let stream: Rc<OnceCell<Class<'js, Self>>> = Rc::new(OnceCell::new());
914
915        // Let iteratorRecord be ? GetIterator(asyncIterable, async).
916        let iterator_record =
917            IteratorRecord::get_iterator(ctx, async_iterable, IteratorKind::Async)?;
918        let iterator = iterator_record.iterator.clone();
919
920        // Let startAlgorithm be an algorithm that returns undefined.
921        let start_algorithm = StartAlgorithm::ReturnUndefined;
922
923        let promise_primordials = PromisePrimordials::get(ctx)?.clone();
924
925        // Let pullAlgorithm be the following steps:
926        let pull_algorithm = {
927            let stream = stream.clone();
928            let promise_primordials = promise_primordials.clone();
929            move |ctx: Ctx<'js>, controller: ReadableStreamControllerClass<'js>| {
930                // Let nextResult be IteratorNext(iteratorRecord).
931                let next_result: Result<Object<'js>> = iterator_record.iterator_next(&ctx, None);
932                let next_promise = match next_result {
933                    // If nextResult is an abrupt completion, return a promise rejected with nextResult.[[Value]].
934                    Err(Error::Exception) => {
935                        return promise_rejected_catch(&ctx, &promise_primordials);
936                    },
937                    Err(err) => return Err(err),
938                    // Let nextPromise be a promise resolved with nextResult.[[Value]].
939                    Ok(next_result) => promise_resolved_with(
940                        &ctx,
941                        &promise_primordials,
942                        Ok(next_result.into_inner()),
943                    )?,
944                };
945
946                // Return the result of reacting to nextPromise with the following fulfillment steps, given iterResult:
947                upon_promise_fulfilment(ctx, next_promise, {
948                    let stream = stream.clone();
949                    move |ctx, iter_result: Value<'js>| {
950                        let iter_result = match iter_result.into_object() {
951                            // If Type(iterResult) is not Object, throw a TypeError.
952                            None => {
953                                return Err(Exception::throw_type(&ctx, "The promise returned by the iterator.next() method must fulfill with an object"));
954                            },
955                            Some(iter_result) => iter_result,
956                        };
957
958                        // Let done be ? IteratorComplete(iterResult).
959                        let done = IteratorRecord::iterator_complete(&iter_result)?;
960
961                        let stream = OwnedBorrowMut::from_class(stream.get().cloned().expect("ReadableStreamFromIterable pull steps called with uninitialised stream"));
962                        let controller = match controller {
963                        ReadableStreamControllerClass::ReadableStreamDefaultController(c) => OwnedBorrowMut::from_class(c),
964                        _ => panic!("ReadableStreamFromIterable pull steps called without default controller")
965                    };
966
967                        let objects = ReadableStreamObjects::new_default(stream, controller);
968
969                        // If done is true:
970                        if done {
971                            // Perform ! ReadableStreamDefaultControllerClose(stream.[[controller]]).
972                            ReadableStreamDefaultController::readable_stream_default_controller_close(ctx.clone(), objects)?;
973                        } else {
974                            // Let value be ? IteratorValue(iterResult).
975                            let value = IteratorRecord::iterator_value(&iter_result)?;
976
977                            // Perform ! ReadableStreamDefaultControllerEnqueue(stream.[[controller]], value).
978                            ReadableStreamDefaultController::readable_stream_default_controller_enqueue(ctx.clone(), objects, value)?;
979                        }
980
981                        Ok(())
982                    }
983                })
984            }
985        };
986
987        // Let cancelAlgorithm be the following steps, given reason:
988        let cancel_algorithm = {
989            let ctx = ctx.clone();
990            let promise_primordials = promise_primordials.clone();
991            move |reason: Value<'js>| {
992                // Let iterator be iteratorRecord.[[Iterator]].
993
994                // Let returnMethod be GetMethod(iterator, "return").
995                let return_method_val: Value<'js> = match iterator.get(PredefinedAtom::Return) {
996                    Err(Error::Exception) => {
997                        return promise_rejected_catch(&ctx, &promise_primordials);
998                    },
999                    Err(err) => return Err(err),
1000                    Ok(val) => val,
1001                };
1002
1003                let return_method: Function<'js> =
1004                    if return_method_val.is_undefined() || return_method_val.is_null() {
1005                        // If returnMethod.[[Value]] is undefined, return a promise resolved with undefined.
1006                        return Ok(promise_primordials.promise_resolved_with_undefined.clone());
1007                    } else if let Some(func) = return_method_val.as_function() {
1008                        func.clone()
1009                    } else {
1010                        // returnMethod is not callable — reject with TypeError
1011                        let _ = Exception::throw_type(&ctx, "return is not a function");
1012                        return promise_rejected_catch(&ctx, &promise_primordials);
1013                    };
1014
1015                // Let returnResult be Call(returnMethod.[[Value]], iterator, « reason »).
1016                let return_result: Result<Value<'js>> =
1017                    return_method.call((This(iterator), reason));
1018
1019                let return_result = match return_result {
1020                    // If returnResult is an abrupt completion, return a promise rejected with returnResult.[[Value]].
1021                    Err(Error::Exception) => {
1022                        return promise_rejected_catch(&ctx, &promise_primordials);
1023                    },
1024                    Err(err) => return Err(err),
1025                    Ok(return_result) => return_result,
1026                };
1027
1028                // Let returnPromise be a promise resolved with returnResult.[[Value]].
1029                let return_promise =
1030                    promise_resolved_with(&ctx, &promise_primordials, Ok(return_result))?;
1031
1032                // Return the result of reacting to returnPromise with the following fulfillment steps, given iterResult:
1033                upon_promise_fulfilment(
1034                    ctx,
1035                    return_promise,
1036                    move |ctx: Ctx<'js>, iter_result: Value<'js>| {
1037                        // If Type(iterResult) is not Object, throw a TypeError.
1038                        if !iter_result.is_object() {
1039                            return Err(Exception::throw_type(&ctx, "The promise returned by the iterator.next() method must fulfill with an object"));
1040                        }
1041                        // Return undefined.
1042                        Ok(rquickjs::Undefined)
1043                    },
1044                )
1045            }
1046        };
1047
1048        let objects_class = ReadableStream::create_readable_stream(
1049            ctx.clone(),
1050            start_algorithm,
1051            PullAlgorithm::from_fn(pull_algorithm),
1052            CancelAlgorithm::from_fn(cancel_algorithm),
1053            Some(0.0),
1054            None,
1055        )?;
1056        _ = stream.set(objects_class.stream.clone());
1057        Ok(objects_class.stream)
1058    }
1059
1060    pub(super) fn reader_mut(&mut self) -> Option<ReadableStreamReaderOwned<'js>> {
1061        self.reader
1062            .clone()
1063            .map(ReadableStreamReaderOwned::from_class)
1064    }
1065}
1066
1067// enum ReadableStreamType { "bytes", "owning" };
1068enum ReadableStreamType {
1069    Bytes,
1070    Owning,
1071}
1072
1073impl<'js> FromJs<'js> for ReadableStreamType {
1074    fn from_js(ctx: &Ctx<'js>, value: Value<'js>) -> Result<Self> {
1075        let typ = value.type_of();
1076
1077        match Coerced::<String>::from_js(ctx, value)?.as_str() {
1078            "bytes" => Ok(Self::Bytes),
1079            "owning" => Ok(Self::Owning),
1080            _ => Err(Error::new_from_js(typ.as_str(), "ReadableStreamType")),
1081        }
1082    }
1083}
1084
1085struct ReadableStreamGetReaderOptions {
1086    mode: Option<ReadableStreamReaderMode>,
1087}
1088
1089impl<'js> FromJs<'js> for ReadableStreamGetReaderOptions {
1090    fn from_js(_ctx: &Ctx<'js>, value: Value<'js>) -> Result<Self> {
1091        let ty_name = value.type_name();
1092        let obj = value
1093            .as_object()
1094            .ok_or(Error::new_from_js(ty_name, "Object"))?;
1095
1096        let mode = obj.get_value_or_undefined::<_, ReadableStreamReaderMode>("mode")?;
1097
1098        Ok(Self { mode })
1099    }
1100}
1101
1102// enum ReadableStreamReaderMode { "byob" };
1103enum ReadableStreamReaderMode {
1104    Byob,
1105}
1106
1107impl<'js> FromJs<'js> for ReadableStreamReaderMode {
1108    fn from_js(ctx: &Ctx<'js>, value: Value<'js>) -> Result<Self> {
1109        let typ = value.type_of();
1110
1111        match Coerced::<String>::from_js(ctx, value)?.as_str() {
1112            "byob" => Ok(Self::Byob),
1113            _ => Err(Error::new_from_js(typ.as_str(), "ReadableStreamReaderMode")),
1114        }
1115    }
1116}