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
58pub 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
69pub 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 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 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 stream_ref.state = ReadableStreamState::Closed;
114 Some(chunks)
115}
116
117pub 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 #[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 let underlying_source = Null(underlying_source.0);
193
194 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 state: ReadableStreamState::Readable,
208 reader: None,
210 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 Some(ReadableStreamType::Bytes) => {
228 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 high_water_mark =
241 QueuingStrategy::extract_high_water_mark(&ctx, queuing_strategy, 0.0)?;
242
243 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 None | Some(ReadableStreamType::Owning) => {
255 let is_owning_type = matches!(
256 underlying_source_dict.r#type,
257 Some(ReadableStreamType::Owning)
258 );
259 let size_algorithm =
261 QueuingStrategy::extract_size_algorithm(queuing_strategy.as_ref());
262
263 let high_water_mark =
265 QueuingStrategy::extract_high_water_mark(&ctx, queuing_strategy, 1.0)?;
266
267 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 #[qjs(static)]
285 fn from(ctx: Ctx<'js>, async_iterable: Value<'js>) -> Result<Class<'js, Self>> {
286 Self::readable_stream_from_iterable(&ctx, async_iterable)
288 }
289
290 #[qjs(get)]
292 fn locked(&self) -> bool {
293 self.is_readable_stream_locked()
295 }
296
297 #[qjs(get)]
299 fn disturbed(&self) -> bool {
300 self.disturbed
301 }
302
303 fn cancel(
305 ctx: Ctx<'js>,
306 stream: This<OwnedBorrowMut<'js, Self>>,
307 reason: Opt<Value<'js>>,
308 ) -> Result<Promise<'js>> {
309 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 fn get_reader(
327 ctx: Ctx<'js>,
328 stream: This<OwnedBorrowMut<'js, Self>>,
329 options: Opt<Option<ReadableStreamGetReaderOptions>>,
330 ) -> Result<ReadableStreamReaderClass<'js>> {
331 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 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 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 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 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 options = options.0.unwrap_or_default();
384
385 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 let () = promise
398 .catch()?
399 .call((This(promise.clone()), Function::new(ctx, || {})))?;
400
401 Ok(readable_class)
403 }
404
405 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 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 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 options = options.unwrap_or_default();
448
449 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 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 (stream, reader) = ReadableStreamReaderClass::acquire_readable_stream_default_reader(
488 ctx.clone(),
489 stream.0,
490 )?;
491
492 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 mut objects: ReadableStreamObjects<'js, C, R>,
521 e: Value<'js>,
522 ) -> Result<ReadableStreamObjects<'js, C, R>> {
523 objects.stream.state = ReadableStreamState::Errored(e.clone());
526
527 objects = objects.with_reader(
528 |mut objects| {
530 objects.reader
532 .generic
533 .closed_promise
534 .reject(e.clone())?;
535
536 objects.reader.generic.closed_promise.set_is_handled()?;
538
539 objects = ReadableStreamDefaultReader::readable_stream_default_reader_error_read_requests(
541 objects, e.clone(),
542 )?;
543 Ok(objects)
544 },
545 |mut objects| {
547 objects.reader
549 .generic
550 .closed_promise
551 .reject(e.clone())?;
552
553 objects.reader.generic.closed_promise.set_is_handled()?;
555
556 objects = ReadableStreamBYOBReader::readable_stream_byob_reader_error_read_into_requests(
558 objects, e.clone(),
559 )?;
560
561 Ok(objects)
562 },
563 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 mut objects: ReadableStreamDefaultReaderObjects<'js, C>,
585 chunk: Value<'js>,
586 done: bool,
587 ) -> Result<ReadableStreamDefaultReaderObjects<'js, C>> {
588 let read_request = objects
591 .reader
592 .read_requests
593 .pop_front()
594 .expect("ReadableStreamFulfillReadRequest called with empty readRequests");
595
596 if done {
597 read_request.close_steps_typed(ctx, objects)
599 } else {
600 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 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 read_into_request.close_steps(objects, chunk.into_js(ctx)?)
622 } else {
623 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 mut objects: ReadableStreamObjects<'js, C, R>,
635 ) -> Result<ReadableStreamObjects<'js, C, R>> {
636 objects.stream.state = ReadableStreamState::Closed;
638
639 objects.with_reader(
640 |mut objects| {
641 objects.reader.generic.closed_promise.resolve_undefined()?;
643
644 let read_requests = objects.reader.read_requests.split_off(0);
648
649 for read_request in read_requests {
651 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 Ok,
664 )
665 }
666
667 pub fn is_readable_stream_locked(&self) -> bool {
668 if self.reader.is_none() {
670 return false;
671 }
672 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 objects.stream.disturbed = true;
694
695 match objects.stream.state {
696 ReadableStreamState::Closed => Ok((
698 promise_resolved_with(
700 &ctx,
701 &objects.stream.promise_primordials,
702 Ok(Value::new_undefined(ctx.clone())),
703 )?,
704 objects,
705 )),
706 ReadableStreamState::Errored(ref stored_error) => Ok((
708 promise_rejected_with(&objects.stream.promise_primordials, stored_error.clone())?,
709 objects,
710 )),
711 ReadableStreamState::Readable => {
712 objects = ReadableStream::readable_stream_close(ctx.clone(), objects)?;
714 objects = objects.with_reader(
718 Ok,
719 |mut objects| {
720 let read_into_requests = objects.reader.read_into_requests.split_off(0);
723 for read_into_request in read_into_requests {
725 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 (source_cancel_promise, objects) = C::cancel_steps(&ctx, objects, reason)?;
737
738 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 reader.read_into_requests.push_back(Box::new(read_request))
754 }
755
756 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 let high_water_mark = high_water_mark.unwrap_or(1.0);
769
770 let size_algorithm = size_algorithm.unwrap_or(SizeAlgorithm::AlwaysOne);
772
773 let base_primordials = BasePrimordials::get(&ctx)?;
774
775 let stream_class = Class::instance(
777 ctx.clone(),
778 Self {
779 state: ReadableStreamState::Readable,
781 reader: None,
783 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 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, )?;
808
809 Ok(ReadableStreamClassObjects {
811 stream: stream_class,
812 controller: controller_class,
813 reader: UndefinedReader,
814 })
815 }
816
817 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 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 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 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_class = Class::instance(
874 ctx.clone(),
875 Self {
876 state: ReadableStreamState::Readable,
878 reader: None,
880 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 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 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 iterator_record =
917 IteratorRecord::get_iterator(ctx, async_iterable, IteratorKind::Async)?;
918 let iterator = iterator_record.iterator.clone();
919
920 let start_algorithm = StartAlgorithm::ReturnUndefined;
922
923 let promise_primordials = PromisePrimordials::get(ctx)?.clone();
924
925 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 next_result: Result<Object<'js>> = iterator_record.iterator_next(&ctx, None);
932 let next_promise = match next_result {
933 Err(Error::Exception) => {
935 return promise_rejected_catch(&ctx, &promise_primordials);
936 },
937 Err(err) => return Err(err),
938 Ok(next_result) => promise_resolved_with(
940 &ctx,
941 &promise_primordials,
942 Ok(next_result.into_inner()),
943 )?,
944 };
945
946 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 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 = 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 {
971 ReadableStreamDefaultController::readable_stream_default_controller_close(ctx.clone(), objects)?;
973 } else {
974 let value = IteratorRecord::iterator_value(&iter_result)?;
976
977 ReadableStreamDefaultController::readable_stream_default_controller_enqueue(ctx.clone(), objects, value)?;
979 }
980
981 Ok(())
982 }
983 })
984 }
985 };
986
987 let cancel_algorithm = {
989 let ctx = ctx.clone();
990 let promise_primordials = promise_primordials.clone();
991 move |reason: Value<'js>| {
992 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 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 let _ = Exception::throw_type(&ctx, "return is not a function");
1012 return promise_rejected_catch(&ctx, &promise_primordials);
1013 };
1014
1015 let return_result: Result<Value<'js>> =
1017 return_method.call((This(iterator), reason));
1018
1019 let return_result = match return_result {
1020 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 return_promise =
1030 promise_resolved_with(&ctx, &promise_primordials, Ok(return_result))?;
1031
1032 upon_promise_fulfilment(
1034 ctx,
1035 return_promise,
1036 move |ctx: Ctx<'js>, iter_result: Value<'js>| {
1037 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 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
1067enum 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
1102enum 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}