1use std::collections::VecDeque;
2
3use crate::utils::{
4 error_messages::ERROR_MSG_ARRAY_BUFFER_DETACHED,
5 option::{Null, Undefined},
6 primordials::{BasePrimordials, Primordial},
7 result::ResultExt,
8};
9use rquickjs::{
10 class::{OwnedBorrow, OwnedBorrowMut, Trace, Tracer},
11 function::Constructor,
12 methods,
13 prelude::{Opt, This},
14 ArrayBuffer, Class, Ctx, Error, Exception, Function, IntoJs, JsLifetime, Object, Promise,
15 Result, TypedArray, Value,
16};
17
18use crate::stream_web::{
19 readable::{
20 byob_reader::{ArrayConstructorPrimordials, ReadableStreamReadIntoRequest, ViewBytes},
21 controller::{
22 ReadableStreamController, ReadableStreamControllerClass, ReadableStreamControllerOwned,
23 },
24 default_controller::ReadableStreamDefaultControllerOwned,
25 default_reader::ReadableStreamReadRequest,
26 objects::{
27 ReadableByteStreamObjects, ReadableStreamBYOBObjects, ReadableStreamClassObjects,
28 ReadableStreamDefaultReaderObjects, ReadableStreamObjects,
29 },
30 reader::ReadableStreamReader,
31 stream::{
32 algorithms::{CancelAlgorithm, PullAlgorithm, StartAlgorithm},
33 source::UnderlyingSource,
34 ReadableStream, ReadableStreamClass, ReadableStreamOwned, ReadableStreamState,
35 },
36 },
37 utils::{
38 class_from_owned_borrow_mut,
39 promise::{promise_resolved_with, upon_promise},
40 UnwrapOrUndefined,
41 },
42};
43
44#[derive(JsLifetime)]
45#[rquickjs::class]
46pub struct ReadableByteStreamController<'js> {
47 auto_allocate_chunk_size: Option<usize>,
48 byob_request: Option<Class<'js, ReadableStreamBYOBRequest<'js>>>,
49 cancel_algorithm: Option<CancelAlgorithm<'js>>,
50 close_requested: bool,
51 pull_again: bool,
52 pull_algorithm: Option<PullAlgorithm<'js>>,
53 pulling: bool,
54 pub(super) pending_pull_intos: VecDeque<PullIntoDescriptor<'js>>,
55 queue: VecDeque<ReadableByteStreamQueueEntry<'js>>,
56 queue_total_size: usize,
57 started: bool,
58 strategy_hwm: f64,
59 pub(super) stream: ReadableStreamClass<'js>,
60
61 pub(super) array_constructor_primordials: ArrayConstructorPrimordials<'js>,
62 constructor_array_buffer: Constructor<'js>,
63 pub(super) function_array_buffer_is_view: Function<'js>,
64}
65
66impl<'js> Trace<'js> for ReadableByteStreamController<'js> {
67 fn trace<'a>(&self, tracer: Tracer<'a, 'js>) {
68 self.auto_allocate_chunk_size.trace(tracer);
69 self.byob_request.trace(tracer);
70 self.cancel_algorithm.trace(tracer);
71 self.pull_algorithm.trace(tracer);
72 self.pending_pull_intos.trace(tracer);
73 self.queue.trace(tracer);
74 self.queue_total_size.trace(tracer);
75 self.started.trace(tracer);
76 self.strategy_hwm.trace(tracer);
77 self.stream.trace(tracer);
78 self.array_constructor_primordials.trace(tracer);
79 self.constructor_array_buffer.trace(tracer);
80 self.function_array_buffer_is_view.trace(tracer);
81 }
82}
83
84pub type ReadableByteStreamControllerClass<'js> = Class<'js, ReadableByteStreamController<'js>>;
85pub(crate) type ReadableByteStreamControllerOwned<'js> =
86 OwnedBorrowMut<'js, ReadableByteStreamController<'js>>;
87
88impl<'js> ReadableByteStreamController<'js> {
89 pub(super) fn set_up_readable_byte_stream_controller_from_underlying_source(
91 ctx: &Ctx<'js>,
92 stream: ReadableStreamOwned<'js>,
93 underlying_source: Null<Undefined<Object<'js>>>,
94 underlying_source_dict: UnderlyingSource<'js>,
95 high_water_mark: f64,
96 ) -> Result<()> {
97 let (start_algorithm, pull_algorithm, cancel_algorithm, auto_allocate_chunk_size) = (
98 underlying_source_dict
101 .start
102 .map(|f| StartAlgorithm::Function {
103 f,
104 underlying_source: underlying_source.clone(),
105 })
106 .unwrap_or(StartAlgorithm::ReturnUndefined),
107 underlying_source_dict
110 .pull
111 .map(|f| PullAlgorithm::Function {
112 f,
113 underlying_source: underlying_source.clone(),
114 })
115 .unwrap_or(PullAlgorithm::ReturnPromiseUndefined),
116 underlying_source_dict
119 .cancel
120 .map(|f| CancelAlgorithm::Function {
121 f,
122 underlying_source,
123 })
124 .unwrap_or(CancelAlgorithm::ReturnPromiseUndefined),
125 underlying_source_dict.auto_allocate_chunk_size,
127 );
128
129 if auto_allocate_chunk_size == Some(0) {
131 return Err(Exception::throw_type(
132 ctx,
133 "autoAllocateChunkSize must be greater than 0",
134 ));
135 }
136
137 Self::set_up_readable_byte_stream_controller(
138 ctx.clone(),
139 stream,
140 start_algorithm,
141 pull_algorithm,
142 cancel_algorithm,
143 high_water_mark,
144 auto_allocate_chunk_size,
145 )?;
146
147 Ok(())
148 }
149
150 pub(super) fn set_up_readable_byte_stream_controller(
151 ctx: Ctx<'js>,
152 stream: ReadableStreamOwned<'js>,
153 start_algorithm: StartAlgorithm<'js>,
154 pull_algorithm: PullAlgorithm<'js>,
155 cancel_algorithm: CancelAlgorithm<'js>,
156 high_water_mark: f64,
157 auto_allocate_chunk_size: Option<usize>,
158 ) -> Result<Class<'js, Self>> {
159 let (stream_class, mut stream) = class_from_owned_borrow_mut(stream);
160
161 let array_constructor_primordials = ArrayConstructorPrimordials::get(&ctx)?.clone();
162 let BasePrimordials {
163 constructor_array_buffer,
164 function_array_buffer_is_view,
165 ..
166 } = &*BasePrimordials::get(&ctx)?;
167
168 let controller = Self {
169 stream: stream_class,
171
172 pull_again: false,
174 pulling: false,
175
176 byob_request: None,
178
179 queue: VecDeque::new(),
181 queue_total_size: 0,
182
183 close_requested: false,
185 started: false,
186
187 strategy_hwm: high_water_mark,
189
190 pull_algorithm: Some(pull_algorithm),
192 cancel_algorithm: Some(cancel_algorithm),
193
194 auto_allocate_chunk_size,
196
197 pending_pull_intos: VecDeque::new(),
198
199 array_constructor_primordials,
200 constructor_array_buffer: constructor_array_buffer.clone(),
201 function_array_buffer_is_view: function_array_buffer_is_view.clone(),
202 };
203
204 let controller_class = Class::instance(ctx.clone(), controller)?;
205
206 stream.controller =
208 ReadableStreamControllerClass::ReadableStreamByteController(controller_class.clone());
209
210 let objects =
211 ReadableStreamObjects::new_byte(stream, OwnedBorrowMut::from_class(controller_class))
212 .refresh_reader();
213
214 let promise_primordials = objects.stream.promise_primordials.clone();
215
216 let (start_result, objects_class) =
218 Self::start_algorithm(ctx.clone(), objects, start_algorithm)?;
219
220 let start_promise = promise_resolved_with(&ctx, &promise_primordials, Ok(start_result))?;
222
223 let _ = upon_promise::<Value<'js>, _>(ctx.clone(), start_promise, {
224 let objects_class = objects_class.clone();
225 move |ctx, result| {
226 let mut objects =
227 ReadableStreamObjects::from_class_no_reader(objects_class).refresh_reader();
228 match result {
229 Ok(_) => {
231 objects.controller.started = true;
233 Self::readable_byte_stream_controller_call_pull_if_needed(ctx, objects)?;
235 Ok(())
236 },
237 Err(r) => {
239 Self::readable_byte_stream_controller_error(objects, r)?;
241 Ok(())
242 },
243 }
244 }
245 })?;
246
247 Ok(objects_class.controller)
248 }
249
250 fn readable_byte_stream_controller_call_pull_if_needed<R: ReadableStreamReader<'js>>(
251 ctx: Ctx<'js>,
252 objects: ReadableByteStreamObjects<'js, R>,
253 ) -> Result<ReadableByteStreamObjects<'js, R>> {
254 let (should_pull, mut objects) =
256 Self::readable_byte_stream_controller_should_call_pull(objects);
257
258 if !should_pull {
260 return Ok(objects);
261 }
262
263 if objects.controller.pulling {
265 objects.controller.pull_again = true;
267
268 return Ok(objects);
270 }
271
272 objects.controller.pulling = true;
274
275 let (pull_promise, objects_class) = Self::pull_algorithm(ctx.clone(), objects)?;
277
278 upon_promise::<Value<'js>, ()>(ctx, pull_promise, {
279 let objects_class = objects_class.clone();
280 move |ctx, result| {
281 let mut objects =
282 ReadableStreamObjects::from_class_no_reader(objects_class).refresh_reader();
283 match result {
284 Ok(_) => {
286 objects.controller.pulling = false;
288 if objects.controller.pull_again {
290 objects.controller.pull_again = false;
292 Self::readable_byte_stream_controller_call_pull_if_needed(
294 ctx, objects,
295 )?;
296 };
297 Ok(())
298 },
299 Err(e) => {
301 Self::readable_byte_stream_controller_error(objects, e)?;
303 Ok(())
304 },
305 }
306 }
307 })?;
308
309 Ok(ReadableStreamObjects::from_class(objects_class))
310 }
311
312 fn readable_byte_stream_controller_should_call_pull<R: ReadableStreamReader<'js>>(
313 mut objects: ReadableByteStreamObjects<'js, R>,
314 ) -> (bool, ReadableByteStreamObjects<'js, R>) {
315 match objects.stream.state {
317 ReadableStreamState::Readable => {},
318 _ => return (false, objects),
320 }
321
322 if objects.controller.close_requested {
324 return (false, objects);
325 }
326
327 if !objects.controller.started {
329 return (false, objects);
330 }
331
332 let (mut has_read_requests, mut has_read_into_requests) = (false, false);
333 objects = objects
334 .with_reader(
335 |objects| {
336 if ReadableStream::readable_stream_get_num_read_requests(&objects.reader) > 0 {
338 has_read_requests = true;
339 }
340 Ok(objects)
341 },
342 |objects| {
343 if ReadableStream::readable_stream_get_num_read_into_requests(&objects.reader)
345 > 0
346 {
347 has_read_into_requests = true;
348 }
349 Ok(objects)
350 },
351 Ok,
352 )
353 .unwrap();
354
355 if has_read_requests || has_read_into_requests {
356 return (true, objects);
357 }
358
359 let desired_size = objects
361 .controller
362 .readable_byte_stream_controller_get_desired_size(&objects.stream);
363
364 if desired_size.0.expect("desired_size must not be null") > 0.0 {
366 return (true, objects);
368 }
369
370 (false, objects)
372 }
373
374 pub(super) fn readable_byte_stream_controller_error<R: ReadableStreamReader<'js>>(
375 mut objects: ReadableByteStreamObjects<'js, R>,
377 e: Value<'js>,
378 ) -> Result<ReadableByteStreamObjects<'js, R>> {
379 if !matches!(objects.stream.state, ReadableStreamState::Readable) {
381 return Ok(objects);
382 };
383
384 objects
386 .controller
387 .readable_byte_stream_controller_clear_pending_pull_intos();
388
389 objects.controller.reset_queue();
391
392 objects
394 .controller
395 .readable_byte_stream_controller_clear_algorithms();
396
397 ReadableStream::readable_stream_error(objects, e)
399 }
400
401 fn readable_byte_stream_controller_clear_pending_pull_intos(&mut self) {
402 self.readable_byte_stream_controller_invalidate_byob_request();
404
405 self.pending_pull_intos.clear();
407 }
408
409 fn readable_byte_stream_controller_invalidate_byob_request(&mut self) {
410 let byob_request = match self.byob_request {
411 None => return,
413 Some(ref byob_request) => byob_request.clone(),
414 };
415 let mut byob_request = OwnedBorrowMut::from_class(byob_request);
416 byob_request.controller = None;
417 byob_request.view = None;
418
419 self.byob_request = None;
420 }
421
422 fn readable_byte_stream_controller_clear_algorithms(&mut self) {
423 self.pull_algorithm = None;
424 self.cancel_algorithm = None;
425 }
426
427 pub(super) fn readable_byte_stream_controller_get_byob_request(
428 ctx: Ctx<'js>,
429 controller: OwnedBorrowMut<'js, Self>,
430 ) -> Result<(
431 Null<Class<'js, ReadableStreamBYOBRequest<'js>>>,
432 OwnedBorrowMut<'js, Self>,
433 )> {
434 if controller.byob_request.is_none() && !controller.pending_pull_intos.is_empty() {
436 let first_descriptor = &controller.pending_pull_intos[0];
438
439 let view = ViewBytes::from_value(
441 &ctx,
442 &controller.function_array_buffer_is_view,
443 Some(
444 &controller
445 .array_constructor_primordials
446 .constructor_uint8array
447 .construct((
448 first_descriptor.buffer.clone(),
449 first_descriptor.byte_offset + first_descriptor.bytes_filled,
450 first_descriptor.byte_length - first_descriptor.bytes_filled,
451 ))?,
452 ),
453 )?;
454
455 let (controller_class, mut controller) = class_from_owned_borrow_mut(controller);
456
457 let byob_request = ReadableStreamBYOBRequest {
459 controller: Some(controller_class),
461 view: Some(view),
463 };
464
465 controller.byob_request = Some(Class::instance(ctx, byob_request)?);
467
468 Ok((Null(controller.byob_request.clone()), controller))
469 } else {
470 Ok((Null(controller.byob_request.clone()), controller))
472 }
473 }
474
475 fn readable_byte_stream_controller_get_desired_size(
476 &self,
477 stream: &ReadableStream<'js>,
478 ) -> Null<f64> {
479 match stream.state {
481 ReadableStreamState::Errored(_) => Null(None),
483 ReadableStreamState::Closed => Null(Some(0.0)),
485 _ => Null(Some(self.strategy_hwm - self.queue_total_size as f64)),
487 }
488 }
489
490 fn reset_queue(&mut self) {
491 self.queue.clear();
493 self.queue_total_size = 0;
495 }
496
497 pub(super) fn readable_byte_stream_controller_close<R: ReadableStreamReader<'js>>(
498 ctx: Ctx<'js>,
499 mut objects: ReadableByteStreamObjects<'js, R>,
501 ) -> Result<ReadableByteStreamObjects<'js, R>> {
502 if objects.controller.close_requested
504 || !matches!(objects.stream.state, ReadableStreamState::Readable)
505 {
506 return Ok(objects);
507 }
508
509 if objects.controller.queue_total_size > 0 {
511 objects.controller.close_requested = true;
513 return Ok(objects);
515 }
516
517 if let Some(first_pending_pull_into) = objects.controller.pending_pull_intos.front() {
520 if first_pending_pull_into.bytes_filled % first_pending_pull_into.element_size != 0 {
522 let e: Value = objects
524 .stream
525 .constructor_type_error
526 .call(("Insufficient bytes to fill elements in the given buffer",))?;
527 Self::readable_byte_stream_controller_error(objects, e.clone())?;
528 return Err(ctx.throw(e));
529 }
530 }
531
532 objects
534 .controller
535 .readable_byte_stream_controller_clear_algorithms();
536
537 ReadableStream::readable_stream_close(ctx, objects)
539 }
540
541 pub(super) fn readable_byte_stream_controller_enqueue<R: ReadableStreamReader<'js>>(
542 ctx: &Ctx<'js>,
543 objects: ReadableByteStreamObjects<'js, R>,
545 chunk: ViewBytes<'js>,
546 ) -> Result<ReadableByteStreamObjects<'js, R>> {
547 Self::readable_byte_stream_controller_enqueue_impl(
548 ctx, objects, chunk, false,
549 )
550 }
551
552 pub(super) fn readable_byte_stream_controller_enqueue_borrowed<R: ReadableStreamReader<'js>>(
567 ctx: &Ctx<'js>,
568 objects: ReadableByteStreamObjects<'js, R>,
569 chunk: ViewBytes<'js>,
570 ) -> Result<ReadableByteStreamObjects<'js, R>> {
571 Self::readable_byte_stream_controller_enqueue_impl(
572 ctx, objects, chunk, true,
573 )
574 }
575
576 fn readable_byte_stream_controller_enqueue_impl<R: ReadableStreamReader<'js>>(
577 ctx: &Ctx<'js>,
578 mut objects: ReadableByteStreamObjects<'js, R>,
579 chunk: ViewBytes<'js>,
580 skip_transfer: bool,
581 ) -> Result<ReadableByteStreamObjects<'js, R>> {
582 if objects.controller.close_requested
584 || !matches!(objects.stream.state, ReadableStreamState::Readable)
585 {
586 return Ok(objects);
587 };
588
589 let (buffer, byte_length, byte_offset) = chunk.get_array_buffer()?;
593
594 buffer.as_raw().ok_or(Exception::throw_type(
596 ctx,
597 "chunk's buffer is detached and so cannot be enqueued",
598 ))?;
599
600 let transferred_buffer = if skip_transfer {
606 buffer
607 } else {
608 transfer_array_buffer(buffer)?
609 };
610
611 if !objects.controller.pending_pull_intos.is_empty() {
614 objects.controller.pending_pull_intos[0]
616 .buffer
617 .as_raw()
618 .or_throw_type(
619 ctx,
620 "The BYOB request's buffer has been detached and so cannot be filled with an enqueued chunk",
621 )?;
622
623 objects
625 .controller
626 .readable_byte_stream_controller_invalidate_byob_request();
627
628 objects.controller.pending_pull_intos[0].buffer =
630 transfer_array_buffer(objects.controller.pending_pull_intos[0].buffer.clone())?;
631
632 if let PullIntoDescriptorReaderType::None =
634 objects.controller.pending_pull_intos[0].reader_type
635 {
636 objects = Self::readable_byte_stream_enqueue_detached_pull_into_to_queue(
637 ctx.clone(),
638 objects,
639 0,
640 )?;
641 }
642 }
643
644 objects = objects.with_reader(
645 |mut objects| {
647 objects = Self::readable_byte_stream_controller_process_read_requests_using_queue(
649 objects, ctx,
650 )?;
651
652 if ReadableStream::readable_stream_get_num_read_requests(&objects.reader) == 0 {
654 objects
656 .controller
657 .readable_byte_stream_controller_enqueue_chunk_to_queue(
658 transferred_buffer.clone(),
659 byte_offset,
660 byte_length,
661 )
662 } else {
663 if !objects.controller.pending_pull_intos.is_empty() {
666 objects
668 .controller
669 .readable_byte_stream_controller_shift_pending_pull_into();
670 }
671
672 let transferred_view = ViewBytes::from_value(
674 ctx,
675 &objects.controller.function_array_buffer_is_view,
676 Some(
677 &objects
678 .controller
679 .array_constructor_primordials
680 .constructor_uint8array
681 .construct((
682 transferred_buffer.clone(),
683 byte_offset,
684 byte_length,
685 ))?,
686 ),
687 );
688
689 objects = ReadableStream::readable_stream_fulfill_read_request(
691 ctx,
692 objects,
693 transferred_view.into_js(ctx)?,
694 false,
695 )?;
696 }
697
698 Ok(objects)
699 },
700 |mut objects| {
701 objects
704 .controller
705 .readable_byte_stream_controller_enqueue_chunk_to_queue(
706 transferred_buffer.clone(),
707 byte_offset,
708 byte_length,
709 );
710 Self::readable_byte_stream_controller_process_pull_into_descriptors_using_queue(
713 ctx, objects,
714 )
715 },
716 |mut objects| {
717 objects
720 .controller
721 .readable_byte_stream_controller_enqueue_chunk_to_queue(
722 transferred_buffer.clone(),
723 byte_offset,
724 byte_length,
725 );
726
727 Ok(objects)
728 },
729 )?;
730
731 Self::readable_byte_stream_controller_call_pull_if_needed(ctx.clone(), objects)
733 }
734
735 fn readable_byte_stream_enqueue_detached_pull_into_to_queue<R: ReadableStreamReader<'js>>(
736 ctx: Ctx<'js>,
737 mut objects: ReadableByteStreamObjects<'js, R>,
738 pull_into_descriptor_index: usize,
739 ) -> Result<ReadableByteStreamObjects<'js, R>> {
740 let pull_into_descriptor =
741 &objects.controller.pending_pull_intos[pull_into_descriptor_index];
742 if pull_into_descriptor.bytes_filled > 0 {
744 let buffer = pull_into_descriptor.buffer.clone();
745 let byte_offset = pull_into_descriptor.byte_offset;
746 let bytes_filled = pull_into_descriptor.bytes_filled;
747 objects = Self::readable_byte_stream_controller_enqueue_cloned_chunk_to_queue(
748 ctx,
749 objects,
750 &buffer,
751 byte_offset,
752 bytes_filled,
753 )?;
754 }
755
756 objects
758 .controller
759 .readable_byte_stream_controller_shift_pending_pull_into();
760
761 Ok(objects)
762 }
763
764 fn readable_byte_stream_controller_process_read_requests_using_queue(
765 mut objects: ReadableStreamDefaultReaderObjects<'js, OwnedBorrowMut<'js, Self>>,
766 ctx: &Ctx<'js>,
767 ) -> Result<ReadableStreamDefaultReaderObjects<'js, OwnedBorrowMut<'js, Self>>> {
768 while !objects.reader.read_requests.is_empty() {
770 if objects.controller.queue_total_size == 0 {
772 return Ok(objects);
773 }
774
775 let read_request = objects.reader.read_requests.pop_front().unwrap();
778 objects = Self::readable_byte_stream_controller_fill_read_request_from_queue(
780 ctx,
781 objects,
782 read_request,
783 )?;
784 }
785
786 Ok(objects)
787 }
788
789 fn readable_byte_stream_controller_shift_pending_pull_into(
790 &mut self,
791 ) -> PullIntoDescriptor<'js> {
792 self.readable_byte_stream_controller_invalidate_byob_request();
794 self.pending_pull_intos.pop_front().expect(
798 "ReadableByteStreamControllerShiftPendingPullInto called on empty pendingPullIntos",
799 )
800 }
801
802 fn readable_byte_stream_controller_enqueue_chunk_to_queue(
803 &mut self,
804 buffer: ArrayBuffer<'js>,
805 byte_offset: usize,
806 byte_length: usize,
807 ) {
808 self.queue.push_back(ReadableByteStreamQueueEntry {
810 buffer,
811 byte_offset,
812 byte_length,
813 });
814
815 self.queue_total_size += byte_length;
817 }
818
819 fn readable_byte_stream_controller_process_pull_into_descriptors_using_queue<
820 R: ReadableStreamReader<'js>,
821 >(
822 ctx: &Ctx<'js>,
823 mut objects: ReadableByteStreamObjects<'js, R>,
824 ) -> Result<ReadableByteStreamObjects<'js, R>> {
825 while !objects.controller.pending_pull_intos.is_empty() {
827 if objects.controller.queue_total_size == 0 {
829 return Ok(objects);
830 }
831
832 let mut pull_into_descriptor_ref = PullIntoDescriptorRefMut::Index(0);
834
835 if objects
837 .controller
838 .readable_byte_stream_controller_fill_pull_into_descriptor_from_queue(
839 ctx,
840 &mut pull_into_descriptor_ref,
841 )?
842 {
843 let pull_into_descriptor = objects
845 .controller
846 .readable_byte_stream_controller_shift_pending_pull_into();
847
848 objects = Self::readable_byte_stream_controller_commit_pull_into_descriptor(
850 ctx.clone(),
851 objects,
852 pull_into_descriptor,
853 )?;
854 }
855 }
856 Ok(objects)
857 }
858
859 fn readable_byte_stream_controller_enqueue_cloned_chunk_to_queue<
860 R: ReadableStreamReader<'js>,
861 >(
862 ctx: Ctx<'js>,
863 mut objects: ReadableByteStreamObjects<'js, R>,
864 buffer: &ArrayBuffer<'js>,
865 byte_offset: usize,
866 byte_length: usize,
867 ) -> Result<ReadableByteStreamObjects<'js, R>> {
868 let clone_result = match ArrayBuffer::new_copy(
870 ctx.clone(),
871 &unsafe { buffer.as_bytes() }.expect(
873 "ReadableByteStreamControllerEnqueueClonedChunkToQueue called on detached buffer",
874 )[byte_offset..byte_offset + byte_length],
875 ) {
876 Ok(clone_result) => clone_result,
877 Err(Error::Exception) => {
878 let err = ctx.catch();
879 Self::readable_byte_stream_controller_error(objects, err.clone())?;
880 return Err(ctx.throw(err));
881 },
882 Err(err) => return Err(err),
883 };
884
885 objects
887 .controller
888 .readable_byte_stream_controller_enqueue_chunk_to_queue(clone_result, 0, byte_length);
889
890 Ok(objects)
891 }
892
893 fn readable_byte_stream_controller_fill_read_request_from_queue(
894 ctx: &Ctx<'js>,
895 mut objects: ReadableStreamDefaultReaderObjects<'js, OwnedBorrowMut<'js, Self>>,
896 read_request: impl ReadableStreamReadRequest<'js>,
897 ) -> Result<ReadableStreamDefaultReaderObjects<'js, OwnedBorrowMut<'js, Self>>> {
898 let entry = {
899 let entry = objects.controller.queue.pop_front().expect(
903 "ReadableByteStreamControllerFillReadRequestFromQueue called with empty queue",
904 );
905
906 objects.controller.queue_total_size -= entry.byte_length;
908
909 entry
910 };
911
912 objects = Self::readable_byte_stream_controller_handle_queue_drain(ctx.clone(), objects)?;
914
915 let view: TypedArray<u8> = objects
917 .controller
918 .array_constructor_primordials
919 .constructor_uint8array
920 .construct((entry.buffer, entry.byte_offset, entry.byte_length))?;
921
922 read_request.chunk_steps_typed(objects, view.into_value())
924 }
925
926 fn readable_byte_stream_controller_fill_pull_into_descriptor_from_queue<'a>(
927 &'a mut self,
928 ctx: &Ctx<'js>,
929 pull_into_descriptor_ref: &mut PullIntoDescriptorRefMut<'js, 'a>,
930 ) -> Result<bool> {
931 let (mut total_bytes_to_copy_remaining, ready) = {
932 let pull_into_descriptor = match pull_into_descriptor_ref {
933 PullIntoDescriptorRefMut::Index(i) => &mut self.pending_pull_intos[*i],
934 PullIntoDescriptorRefMut::Owned(r) => r,
935 };
936 let max_bytes_to_copy: usize = std::cmp::min(
938 self.queue_total_size,
939 pull_into_descriptor.byte_length - pull_into_descriptor.bytes_filled,
940 );
941
942 let max_bytes_filled = pull_into_descriptor.bytes_filled + max_bytes_to_copy;
944
945 let mut total_bytes_to_copy_remaining = max_bytes_to_copy;
947
948 let mut ready = false;
950
951 let remainder_bytes = max_bytes_filled % pull_into_descriptor.element_size;
953
954 let max_aligned_bytes = max_bytes_filled - remainder_bytes;
956
957 if max_aligned_bytes >= pull_into_descriptor.minimum_fill {
959 total_bytes_to_copy_remaining =
961 max_aligned_bytes - pull_into_descriptor.bytes_filled;
962 ready = true
964 }
965
966 (total_bytes_to_copy_remaining, ready)
967 };
968
969 while total_bytes_to_copy_remaining > 0 {
972 let bytes_to_copy = {
973 let pull_into_descriptor = match pull_into_descriptor_ref {
974 PullIntoDescriptorRefMut::Index(i) => &mut self.pending_pull_intos[*i],
975 PullIntoDescriptorRefMut::Owned(r) => r,
976 };
977
978 let head_of_queue = self
980 .queue
981 .front_mut()
982 .expect("empty queue with bytes to copy");
983 let bytes_to_copy: usize =
985 std::cmp::min(total_bytes_to_copy_remaining, head_of_queue.byte_length);
986 let dest_start: usize =
988 pull_into_descriptor.byte_offset + pull_into_descriptor.bytes_filled;
989 copy_data_block_bytes(
991 ctx,
992 &pull_into_descriptor.buffer,
993 dest_start,
994 &head_of_queue.buffer,
995 head_of_queue.byte_offset,
996 bytes_to_copy,
997 )?;
998 if head_of_queue.byte_length == bytes_to_copy {
999 self.queue.pop_front();
1002 } else {
1003 head_of_queue.byte_offset += bytes_to_copy;
1006 head_of_queue.byte_length -= bytes_to_copy
1008 }
1009
1010 self.queue_total_size -= bytes_to_copy;
1012
1013 bytes_to_copy
1014 };
1015
1016 self.readable_byte_stream_controller_fill_head_pull_into_descriptor(
1018 bytes_to_copy,
1019 pull_into_descriptor_ref,
1020 );
1021
1022 total_bytes_to_copy_remaining -= bytes_to_copy
1024 }
1025
1026 Ok(ready)
1027 }
1028
1029 fn readable_byte_stream_controller_commit_pull_into_descriptor<R: ReadableStreamReader<'js>>(
1030 ctx: Ctx<'js>,
1031 objects: ReadableByteStreamObjects<'js, R>,
1032 pull_into_descriptor: PullIntoDescriptor<'js>,
1033 ) -> Result<ReadableByteStreamObjects<'js, R>> {
1034 let mut done = false;
1036 if matches!(objects.stream.state, ReadableStreamState::Closed) {
1038 done = true
1040 }
1041
1042 let reader_type = pull_into_descriptor.reader_type;
1043
1044 let filled_view = Self::readable_byte_stream_controller_convert_pull_into_descriptor(
1046 ctx.clone(),
1047 &objects.stream.function_array_buffer_is_view,
1048 pull_into_descriptor,
1049 )?;
1050
1051 if let PullIntoDescriptorReaderType::Default = reader_type {
1052 objects.with_assert_default_reader(|objects| {
1054 ReadableStream::readable_stream_fulfill_read_request(
1056 &ctx,
1057 objects,
1058 filled_view.into_js(&ctx)?,
1059 done,
1060 )
1061 })
1062 } else {
1063 objects.with_assert_byob_reader(|objects| {
1065 ReadableStream::readable_stream_fulfill_read_into_request(
1067 &ctx,
1068 objects,
1069 filled_view,
1070 done,
1071 )
1072 })
1073 }
1074 }
1075
1076 fn readable_byte_stream_controller_handle_queue_drain<R: ReadableStreamReader<'js>>(
1077 ctx: Ctx<'js>,
1078 mut objects: ReadableByteStreamObjects<'js, R>,
1079 ) -> Result<ReadableByteStreamObjects<'js, R>> {
1080 if objects.controller.queue_total_size == 0 && objects.controller.close_requested {
1082 objects
1084 .controller
1085 .readable_byte_stream_controller_clear_algorithms();
1086 ReadableStream::readable_stream_close(ctx, objects)
1088 } else {
1089 Self::readable_byte_stream_controller_call_pull_if_needed(ctx.clone(), objects)
1092 }
1093 }
1094
1095 fn readable_byte_stream_controller_convert_pull_into_descriptor(
1096 ctx: Ctx<'js>,
1097 function_array_buffer_is_view: &Function<'js>,
1098 pull_into_descriptor: PullIntoDescriptor<'js>,
1099 ) -> Result<ViewBytes<'js>> {
1100 let PullIntoDescriptor {
1101 bytes_filled,
1103 element_size,
1105 byte_offset,
1106 buffer,
1107 ..
1108 } = pull_into_descriptor;
1109 let buffer = transfer_array_buffer(buffer);
1111 let view: Object = pull_into_descriptor.view_constructor.construct((
1113 buffer,
1114 byte_offset,
1115 bytes_filled / element_size,
1116 ))?;
1117 ViewBytes::from_object(&ctx, function_array_buffer_is_view, &view)
1118 }
1119
1120 pub(super) fn readable_byte_stream_controller_pull_into(
1121 ctx: &Ctx<'js>,
1122 mut objects: ReadableStreamBYOBObjects<'js>,
1124 view: ViewBytes<'js>,
1125 min: u64,
1126 read_into_request: impl ReadableStreamReadIntoRequest<'js> + 'js,
1127 ) -> Result<ReadableStreamBYOBObjects<'js>> {
1128 let (element_size, ctor) = (
1131 view.element_size(),
1132 objects
1133 .controller
1134 .array_constructor_primordials
1135 .for_view_bytes(&view),
1136 );
1137
1138 let minimum_fill: usize = (min as usize) * element_size;
1140
1141 let (buffer, byte_length, byte_offset) = view.get_array_buffer()?;
1144
1145 let buffer_result = transfer_array_buffer(buffer);
1147 let buffer = match buffer_result {
1148 Err(Error::Exception) => {
1150 objects = read_into_request.error_steps(objects, ctx.catch())?;
1152 return Ok(objects);
1154 },
1155 Err(err) => return Err(err),
1156 Ok(buffer) => buffer,
1158 };
1159
1160 let buffer_byte_length = buffer.len();
1161 let mut pull_into_descriptor = PullIntoDescriptor {
1163 buffer,
1164 buffer_byte_length,
1165 byte_offset,
1166 byte_length,
1167 bytes_filled: 0,
1168 minimum_fill,
1169 element_size,
1170 view_constructor: ctor.clone(),
1171 reader_type: PullIntoDescriptorReaderType::Byob,
1172 };
1173
1174 if !objects.controller.pending_pull_intos.is_empty() {
1176 objects
1178 .controller
1179 .pending_pull_intos
1180 .push_back(pull_into_descriptor);
1181
1182 ReadableStream::readable_stream_add_read_into_request(
1184 &mut objects.reader,
1185 read_into_request,
1186 );
1187
1188 return Ok(objects);
1190 }
1191
1192 if matches!(objects.stream.state, ReadableStreamState::Closed) {
1194 let empty_view: Value<'js> = ctor.construct((
1196 pull_into_descriptor.buffer,
1197 pull_into_descriptor.byte_offset,
1198 0,
1199 ))?;
1200
1201 objects = read_into_request.close_steps(objects, empty_view)?;
1203
1204 return Ok(objects);
1206 }
1207
1208 if objects.controller.queue_total_size > 0 {
1210 if objects
1212 .controller
1213 .readable_byte_stream_controller_fill_pull_into_descriptor_from_queue(
1214 ctx,
1215 &mut PullIntoDescriptorRefMut::Owned(&mut pull_into_descriptor),
1216 )?
1217 {
1218 let filled_view = objects
1220 .controller
1221 .readable_byte_steam_controller_convert_pull_into_descriptor(
1222 pull_into_descriptor,
1223 )?;
1224
1225 objects =
1227 Self::readable_byte_stream_controller_handle_queue_drain(ctx.clone(), objects)?;
1228
1229 return read_into_request.chunk_steps(objects, filled_view);
1232 }
1233
1234 if objects.controller.close_requested {
1236 let e: Value = objects
1238 .stream
1239 .constructor_type_error
1240 .call(("Insufficient bytes to fill elements in the given buffer",))?;
1241
1242 objects = Self::readable_byte_stream_controller_error(objects, e.clone())?;
1244
1245 return read_into_request.error_steps(objects, e);
1248 }
1249 }
1250
1251 objects
1253 .controller
1254 .pending_pull_intos
1255 .push_back(pull_into_descriptor);
1256
1257 ReadableStream::readable_stream_add_read_into_request(
1259 &mut objects.reader,
1260 read_into_request,
1261 );
1262
1263 Self::readable_byte_stream_controller_call_pull_if_needed(ctx.clone(), objects)
1265 }
1266
1267 fn readable_byte_steam_controller_convert_pull_into_descriptor(
1268 &mut self,
1269 pull_into_descriptor: PullIntoDescriptor<'js>,
1270 ) -> Result<Value<'js>> {
1271 let bytes_filled = pull_into_descriptor.bytes_filled;
1273
1274 let element_size = pull_into_descriptor.element_size;
1276
1277 let buffer = transfer_array_buffer(pull_into_descriptor.buffer)?;
1279
1280 pull_into_descriptor.view_constructor.construct((
1282 buffer,
1283 pull_into_descriptor.byte_offset,
1284 bytes_filled / element_size,
1285 ))
1286 }
1287
1288 pub(super) fn readable_byte_stream_controller_respond<R: ReadableStreamReader<'js>>(
1289 ctx: Ctx<'js>,
1290 mut objects: ReadableByteStreamObjects<'js, R>,
1291 bytes_written: usize,
1292 ) -> Result<()> {
1293 let first_descriptor = &mut objects.controller.pending_pull_intos[0];
1295
1296 match objects.stream.state {
1298 ReadableStreamState::Closed => {
1300 if bytes_written != 0 {
1302 return Err(Exception::throw_type(
1303 &ctx,
1304 "bytesWritten must be 0 when calling respond() on a closed stream",
1305 ));
1306 }
1307 },
1308 _ => {
1310 if bytes_written == 0 {
1312 return Err(Exception::throw_type(
1313 &ctx,
1314 "bytesWritten must be greater than 0 when calling respond() on a readable stream",
1315 ));
1316 }
1317
1318 if first_descriptor.bytes_filled + bytes_written > first_descriptor.byte_length {
1320 return Err(Exception::throw_range(&ctx, "bytesWritten out of range'"));
1321 }
1322 },
1323 };
1324
1325 first_descriptor.buffer = transfer_array_buffer(first_descriptor.buffer.clone())?;
1327
1328 Self::readable_byte_stream_controller_respond_internal(ctx, objects, bytes_written)
1330 }
1331
1332 fn readable_byte_stream_controller_respond_internal<R: ReadableStreamReader<'js>>(
1333 ctx: Ctx<'js>,
1334 mut objects: ReadableByteStreamObjects<'js, R>,
1335 bytes_written: usize,
1336 ) -> Result<()> {
1337 let first_descriptor_index = 0;
1339
1340 objects
1342 .controller
1343 .readable_byte_stream_controller_invalidate_byob_request();
1344
1345 match objects.stream.state {
1347 ReadableStreamState::Closed => {
1349 objects = Self::readable_byte_stream_controller_respond_in_closed_state(
1351 ctx.clone(),
1352 objects,
1353 first_descriptor_index,
1354 )?;
1355 },
1356 _ => {
1358 objects = Self::readable_byte_stream_controller_respond_in_readable_state(
1360 ctx.clone(),
1361 objects,
1362 bytes_written,
1363 first_descriptor_index,
1364 )?
1365 },
1366 };
1367
1368 _ = Self::readable_byte_stream_controller_call_pull_if_needed(ctx, objects)?;
1369 Ok(())
1370 }
1371
1372 fn readable_byte_stream_controller_respond_in_closed_state<R: ReadableStreamReader<'js>>(
1373 ctx: Ctx<'js>,
1374 mut objects: ReadableByteStreamObjects<'js, R>,
1376 first_descriptor_index: usize,
1377 ) -> Result<ReadableByteStreamObjects<'js, R>> {
1378 if let PullIntoDescriptorReaderType::None =
1380 objects.controller.pending_pull_intos[first_descriptor_index].reader_type
1381 {
1382 objects
1383 .controller
1384 .readable_byte_stream_controller_shift_pending_pull_into();
1385 }
1386
1387 objects.with_reader(
1389 Ok,
1390 |mut objects| {
1391 while ReadableStream::readable_stream_get_num_read_into_requests(&objects.reader)
1393 > 0
1394 {
1395 let pull_into_descriptor = objects
1397 .controller
1398 .readable_byte_stream_controller_shift_pending_pull_into();
1399
1400 objects = Self::readable_byte_stream_controller_commit_pull_into_descriptor(
1402 ctx.clone(),
1403 objects,
1404 pull_into_descriptor,
1405 )?;
1406 }
1407
1408 Ok(objects)
1409 },
1410 Ok,
1411 )
1412 }
1413
1414 fn readable_byte_stream_controller_respond_in_readable_state<R: ReadableStreamReader<'js>>(
1415 ctx: Ctx<'js>,
1416 mut objects: ReadableByteStreamObjects<'js, R>,
1418 bytes_written: usize,
1419 pull_into_descriptor_index: usize,
1420 ) -> Result<ReadableByteStreamObjects<'js, R>> {
1421 objects
1423 .controller
1424 .readable_byte_stream_controller_fill_head_pull_into_descriptor(
1425 bytes_written,
1426 &mut PullIntoDescriptorRefMut::Index(pull_into_descriptor_index),
1427 );
1428
1429 if let PullIntoDescriptorReaderType::None =
1431 objects.controller.pending_pull_intos[pull_into_descriptor_index].reader_type
1432 {
1433 objects = Self::readable_byte_stream_enqueue_detached_pull_into_to_queue(
1435 ctx.clone(),
1436 objects,
1437 pull_into_descriptor_index,
1438 )?;
1439 return Self::readable_byte_stream_controller_process_pull_into_descriptors_using_queue(
1442 &ctx, objects,
1443 );
1444 }
1445
1446 if objects.controller.pending_pull_intos[pull_into_descriptor_index].bytes_filled
1448 < objects.controller.pending_pull_intos[pull_into_descriptor_index].minimum_fill
1449 {
1450 return Ok(objects);
1451 }
1452
1453 let mut pull_into_descriptor = objects
1455 .controller
1456 .readable_byte_stream_controller_shift_pending_pull_into();
1457
1458 let remainder_size = pull_into_descriptor.bytes_filled % pull_into_descriptor.element_size;
1460
1461 if remainder_size > 0 {
1463 let end = pull_into_descriptor.byte_offset + pull_into_descriptor.bytes_filled;
1465
1466 let buffer = pull_into_descriptor.buffer.clone();
1467
1468 objects = Self::readable_byte_stream_controller_enqueue_cloned_chunk_to_queue(
1470 ctx.clone(),
1471 objects,
1472 &buffer,
1473 end - remainder_size,
1474 remainder_size,
1475 )?;
1476 }
1477
1478 pull_into_descriptor.bytes_filled -= remainder_size;
1480
1481 objects = Self::readable_byte_stream_controller_commit_pull_into_descriptor(
1483 ctx.clone(),
1484 objects,
1485 pull_into_descriptor,
1486 )?;
1487
1488 Self::readable_byte_stream_controller_process_pull_into_descriptors_using_queue(
1490 &ctx, objects,
1491 )
1492 }
1493
1494 pub(super) fn readable_byte_stream_controller_respond_with_new_view<
1495 R: ReadableStreamReader<'js>,
1496 >(
1497 ctx: Ctx<'js>,
1498 mut objects: ReadableByteStreamObjects<'js, R>,
1499 view: ViewBytes<'js>,
1500 ) -> Result<()> {
1501 let first_descriptor_index = 0;
1503
1504 let (buffer, byte_length, byte_offset) = view.get_array_buffer()?;
1505
1506 match objects.stream.state {
1508 ReadableStreamState::Closed => {
1510 if byte_length != 0 {
1512 return Err(Exception::throw_type(&ctx, "The view's length must be 0 when calling respondWithNewView() on a closed stream"));
1513 }
1514 },
1515 _ => {
1517 if byte_length == 0 {
1519 return Err(Exception::throw_type(&ctx, "The view's length must be greater than 0 when calling respondWithNewView() on a readable stream"));
1520 }
1521 },
1522 };
1523
1524 {
1525 let first_descriptor =
1526 &mut objects.controller.pending_pull_intos[first_descriptor_index];
1527
1528 if first_descriptor.byte_offset + first_descriptor.bytes_filled != byte_offset {
1530 return Err(Exception::throw_range(
1531 &ctx,
1532 "The region specified by view does not match byobRequest",
1533 ));
1534 };
1535
1536 if first_descriptor.buffer_byte_length != buffer.len() {
1538 return Err(Exception::throw_range(
1539 &ctx,
1540 "The buffer of view has different capacity than byobRequest",
1541 ));
1542 };
1543
1544 if first_descriptor.bytes_filled + byte_length > first_descriptor.byte_length {
1546 return Err(Exception::throw_range(
1547 &ctx,
1548 "The region specified by view is larger than byobRequest",
1549 ));
1550 }
1551
1552 first_descriptor.buffer = transfer_array_buffer(buffer)?;
1554 }
1555
1556 Self::readable_byte_stream_controller_respond_internal(ctx, objects, byte_length)
1558 }
1559
1560 fn readable_byte_stream_controller_fill_head_pull_into_descriptor<'a>(
1561 &mut self,
1562 size: usize,
1563 pull_into_descriptor_ref: &mut PullIntoDescriptorRefMut<'js, 'a>,
1564 ) {
1565 let pull_into_descriptor = match pull_into_descriptor_ref {
1566 PullIntoDescriptorRefMut::Index(i) => &mut self.pending_pull_intos[*i],
1567 PullIntoDescriptorRefMut::Owned(r) => *r,
1568 };
1569
1570 pull_into_descriptor.bytes_filled += size;
1572 }
1573
1574 fn start_algorithm<R: ReadableStreamReader<'js>>(
1575 ctx: Ctx<'js>,
1576 objects: ReadableByteStreamObjects<'js, R>,
1577 start_algorithm: StartAlgorithm<'js>,
1578 ) -> Result<(
1579 Value<'js>,
1580 ReadableStreamClassObjects<'js, OwnedBorrowMut<'js, Self>, R>,
1581 )> {
1582 let objects_class = objects.into_inner();
1583
1584 Ok((
1585 start_algorithm.call(
1586 ctx,
1587 ReadableStreamControllerClass::ReadableStreamByteController(
1588 objects_class.controller.clone(),
1589 ),
1590 )?,
1591 objects_class,
1592 ))
1593 }
1594
1595 fn pull_algorithm<R: ReadableStreamReader<'js>>(
1596 ctx: Ctx<'js>,
1597 objects: ReadableByteStreamObjects<'js, R>,
1598 ) -> Result<(
1599 Promise<'js>,
1600 ReadableStreamClassObjects<'js, OwnedBorrowMut<'js, Self>, R>,
1601 )> {
1602 let pull_algorithm = objects
1603 .controller
1604 .pull_algorithm
1605 .clone()
1606 .expect("pull algorithm used after ReadableStreamDefaultControllerClearAlgorithms");
1607 let promise_primordials = objects.stream.promise_primordials.clone();
1608 let objects_class = objects.into_inner();
1609
1610 Ok((
1611 pull_algorithm.call(
1612 ctx,
1613 &promise_primordials,
1614 ReadableStreamControllerClass::ReadableStreamByteController(
1615 objects_class.controller.clone(),
1616 ),
1617 )?,
1618 objects_class,
1619 ))
1620 }
1621
1622 fn cancel_algorithm<R: ReadableStreamReader<'js>>(
1623 ctx: Ctx<'js>,
1624 objects: ReadableByteStreamObjects<'js, R>,
1625 reason: Value<'js>,
1626 ) -> Result<(
1627 Promise<'js>,
1628 ReadableStreamClassObjects<'js, OwnedBorrowMut<'js, Self>, R>,
1629 )> {
1630 let cancel_algorithm =
1631 objects.controller.cancel_algorithm.clone().expect(
1632 "cancel algorithm used after ReadableStreamDefaultControllerClearAlgorithms",
1633 );
1634 let promise_primordials = objects.stream.promise_primordials.clone();
1635 let objects_class = objects.into_inner();
1636
1637 Ok((
1638 cancel_algorithm.call(ctx, &promise_primordials, reason)?,
1639 objects_class,
1640 ))
1641 }
1642}
1643
1644#[methods(rename_all = "camelCase")]
1645impl<'js> ReadableByteStreamController<'js> {
1646 #[qjs(constructor)]
1647 fn new(ctx: Ctx<'js>) -> Result<Class<'js, Self>> {
1648 Err(Exception::throw_type(&ctx, "Illegal constructor"))
1649 }
1650
1651 #[qjs(get, rename = "byobRequest")]
1653 fn byob_request_getter(
1654 ctx: Ctx<'js>,
1655 controller: This<Class<'js, Self>>,
1656 ) -> Result<Null<Class<'js, ReadableStreamBYOBRequest<'js>>>> {
1657 if let Ok(owned) = rquickjs::class::OwnedBorrowMut::try_from_class(controller.0.clone()) {
1665 let (request, _) = Self::readable_byte_stream_controller_get_byob_request(ctx, owned)?;
1666 return Ok(request);
1667 }
1668 Ok(Null(None))
1674 }
1675
1676 #[qjs(get)]
1678 fn desired_size(&self) -> Null<f64> {
1679 let stream = OwnedBorrow::from_class(self.stream.clone());
1680 self.readable_byte_stream_controller_get_desired_size(&stream)
1681 }
1682
1683 fn close(ctx: Ctx<'js>, controller: This<OwnedBorrowMut<'js, Self>>) -> Result<()> {
1685 if controller.close_requested {
1687 return Err(Exception::throw_type(&ctx, "close() called more than once"));
1688 }
1689
1690 let objects = ReadableStreamObjects::from_byte_controller(controller.0).refresh_reader();
1691
1692 if !matches!(objects.stream.state, ReadableStreamState::Readable) {
1693 return Err(Exception::throw_type(
1694 &ctx,
1695 "close() called when stream is not readable",
1696 ));
1697 };
1698
1699 Self::readable_byte_stream_controller_close(ctx, objects)?;
1701 Ok(())
1702 }
1703
1704 fn enqueue(
1706 this: This<OwnedBorrowMut<'js, Self>>,
1707 ctx: Ctx<'js>,
1708 chunk: Value<'js>,
1709 ) -> Result<()> {
1710 let chunk = ViewBytes::from_value(&ctx, &this.function_array_buffer_is_view, Some(&chunk))?;
1711
1712 let (array_buffer, byte_length, _) = chunk.get_array_buffer()?;
1713
1714 if byte_length == 0 {
1716 return Err(Exception::throw_type(
1717 &ctx,
1718 "chunk must have non-zero byteLength",
1719 ));
1720 }
1721
1722 if array_buffer.is_empty() {
1724 return Err(Exception::throw_type(
1725 &ctx,
1726 "chunk must have non-zero buffer byteLength",
1727 ));
1728 }
1729
1730 if this.close_requested {
1732 return Err(Exception::throw_type(&ctx, "stream is closed or draining"));
1733 }
1734
1735 let objects = ReadableStreamObjects::from_byte_controller(this.0).refresh_reader();
1736
1737 if !matches!(objects.stream.state, ReadableStreamState::Readable) {
1739 return Err(Exception::throw_type(
1740 &ctx,
1741 "The stream is not in the readable state and cannot be enqueued to",
1742 ));
1743 };
1744
1745 Self::readable_byte_stream_controller_enqueue(&ctx, objects, chunk)?;
1747 Ok(())
1748 }
1749
1750 fn error(
1752 ctx: Ctx<'js>,
1753 controller: This<OwnedBorrowMut<'js, Self>>,
1754 e: Opt<Value<'js>>,
1755 ) -> Result<()> {
1756 let objects = ReadableStreamObjects::from_byte_controller(controller.0).refresh_reader();
1757
1758 Self::readable_byte_stream_controller_error(objects, e.0.unwrap_or_undefined(&ctx))?;
1760 Ok(())
1761 }
1762}
1763
1764impl<'js> ReadableStreamController<'js> for ReadableByteStreamControllerOwned<'js> {
1765 type Class = ReadableByteStreamControllerClass<'js>;
1766
1767 fn with_controller<C, O>(
1768 self,
1769 ctx: C,
1770 _: impl FnOnce(
1771 C,
1772 ReadableStreamDefaultControllerOwned<'js>,
1773 ) -> Result<(O, ReadableStreamDefaultControllerOwned<'js>)>,
1774 byte: impl FnOnce(
1775 C,
1776 ReadableByteStreamControllerOwned<'js>,
1777 ) -> Result<(O, ReadableByteStreamControllerOwned<'js>)>,
1778 ) -> Result<(O, Self)> {
1779 let (ctx, reader) = byte(ctx, self)?;
1780 Ok((ctx, reader))
1781 }
1782
1783 fn into_inner(self) -> Self::Class {
1784 OwnedBorrowMut::into_inner(self)
1785 }
1786
1787 fn from_class(class: Self::Class) -> Self {
1788 OwnedBorrowMut::from_class(class)
1789 }
1790
1791 fn into_erased(self) -> ReadableStreamControllerOwned<'js> {
1792 ReadableStreamControllerOwned::ReadableStreamByteController(self)
1793 }
1794
1795 fn try_from_erased(erased: ReadableStreamControllerOwned<'js>) -> Option<Self> {
1796 match erased {
1797 ReadableStreamControllerOwned::ReadableStreamDefaultController(_) => None,
1798 ReadableStreamControllerOwned::ReadableStreamByteController(r) => Some(r),
1799 }
1800 }
1801
1802 fn pull_steps(
1803 ctx: &Ctx<'js>,
1804 mut objects: ReadableStreamDefaultReaderObjects<'js, Self>,
1805 read_request: impl ReadableStreamReadRequest<'js> + 'js,
1806 ) -> Result<ReadableStreamDefaultReaderObjects<'js, Self>> {
1807 if objects.controller.queue_total_size > 0 {
1809 return ReadableByteStreamController::readable_byte_stream_controller_fill_read_request_from_queue(
1812 ctx,
1813 objects,
1814 read_request,
1815 );
1816 }
1817
1818 let auto_allocate_chunk_size = objects.controller.auto_allocate_chunk_size;
1820
1821 if let Some(auto_allocate_chunk_size) = auto_allocate_chunk_size {
1823 let buffer: ArrayBuffer = match objects
1825 .controller
1826 .constructor_array_buffer
1827 .construct((auto_allocate_chunk_size,))
1828 {
1829 Err(Error::Exception) => {
1831 return read_request.error_steps_typed(objects, ctx.catch());
1833 },
1834 Err(err) => return Err(err),
1835 Ok(buffer) => buffer,
1836 };
1837
1838 let pull_into_descriptor = PullIntoDescriptor {
1840 buffer,
1841 buffer_byte_length: auto_allocate_chunk_size,
1842 byte_offset: 0,
1843 byte_length: auto_allocate_chunk_size,
1844 bytes_filled: 0,
1845 minimum_fill: 1,
1846 element_size: 1,
1847 view_constructor: objects
1848 .controller
1849 .array_constructor_primordials
1850 .constructor_uint8array
1851 .clone(),
1852 reader_type: PullIntoDescriptorReaderType::Default,
1853 };
1854
1855 objects
1857 .controller
1858 .pending_pull_intos
1859 .push_back(pull_into_descriptor);
1860 }
1861
1862 objects
1864 .stream
1865 .readable_stream_add_read_request(&mut objects.reader, read_request);
1866
1867 ReadableByteStreamController::readable_byte_stream_controller_call_pull_if_needed(
1869 ctx.clone(),
1870 objects,
1871 )
1872 }
1873
1874 fn cancel_steps<R: ReadableStreamReader<'js>>(
1875 ctx: &Ctx<'js>,
1876 mut objects: ReadableStreamObjects<'js, Self, R>,
1877 reason: Value<'js>,
1878 ) -> Result<(Promise<'js>, ReadableStreamObjects<'js, Self, R>)> {
1879 objects
1881 .controller
1882 .readable_byte_stream_controller_clear_pending_pull_intos();
1883
1884 objects.controller.reset_queue();
1886
1887 let (result, objects_class) =
1889 ReadableByteStreamController::cancel_algorithm(ctx.clone(), objects, reason)?;
1890
1891 objects = ReadableStreamObjects::from_class(objects_class);
1892
1893 objects
1895 .controller
1896 .readable_byte_stream_controller_clear_algorithms();
1897
1898 Ok((result, objects))
1900 }
1901
1902 fn release_steps(&mut self) {
1903 if !self.pending_pull_intos.is_empty() {
1905 let first_pending_pull_into = &mut self.pending_pull_intos[0];
1907
1908 first_pending_pull_into.reader_type = PullIntoDescriptorReaderType::None;
1910
1911 _ = self.pending_pull_intos.split_off(1);
1913 }
1914 }
1915}
1916
1917#[derive(JsLifetime, Trace, Clone)]
1918#[rquickjs::class]
1919pub(crate) struct ReadableStreamBYOBRequest<'js> {
1920 pub(super) view: Option<ViewBytes<'js>>,
1921 controller: Option<ReadableByteStreamControllerClass<'js>>,
1922}
1923
1924#[methods(rename_all = "camelCase")]
1925impl<'js> ReadableStreamBYOBRequest<'js> {
1926 #[qjs(constructor)]
1927 fn new(ctx: Ctx<'js>) -> Result<Class<'js, Self>> {
1928 Err(Exception::throw_type(&ctx, "Illegal constructor"))
1929 }
1930
1931 #[qjs(get)]
1932 fn view(&self) -> Null<ViewBytes<'js>> {
1933 Null(self.view.clone())
1934 }
1935
1936 fn respond(
1937 ctx: Ctx<'js>,
1938 byob_request: This<OwnedBorrowMut<'js, Self>>,
1939 bytes_written: usize,
1940 ) -> Result<()> {
1941 let (controller, view) = match (&byob_request.controller, &byob_request.view) {
1943 (Some(controller), Some(view)) => (controller.clone(), view),
1944 _ => {
1945 return Err(Exception::throw_type(
1946 &ctx,
1947 "This BYOB request has been invalidated",
1948 ));
1949 },
1950 };
1951 let (buffer, _, _) = view.get_array_buffer()?;
1952 drop(byob_request);
1953
1954 if unsafe { buffer.as_bytes() }.is_none() {
1957 return Err(Exception::throw_type(
1958 &ctx,
1959 "The BYOB request's buffer has been detached and so cannot be used as a response",
1960 ));
1961 }
1962
1963 let objects =
1964 ReadableStreamObjects::from_byte_controller(OwnedBorrowMut::from_class(controller))
1965 .refresh_reader();
1966
1967 ReadableByteStreamController::readable_byte_stream_controller_respond(
1969 ctx,
1970 objects,
1971 bytes_written,
1972 )
1973 }
1974
1975 fn respond_with_new_view(
1976 ctx: Ctx<'js>,
1977 byob_request: This<OwnedBorrowMut<'js, Self>>,
1978 view: Opt<Value<'js>>,
1979 ) -> Result<()> {
1980 let controller = match &byob_request.controller {
1982 Some(controller) => controller.clone(),
1983 _ => {
1984 return Err(Exception::throw_type(
1985 &ctx,
1986 "This BYOB request has been invalidated",
1987 ));
1988 },
1989 };
1990 drop(byob_request);
1991
1992 let controller = OwnedBorrowMut::from_class(controller);
1993
1994 let view = ViewBytes::from_value(
1995 &ctx,
1996 &controller.function_array_buffer_is_view,
1997 view.0.as_ref(),
1998 )?;
1999
2000 let (buffer, _, _) = view.get_array_buffer()?;
2001
2002 if unsafe { buffer.as_bytes() }.is_none() {
2005 return Err(Exception::throw_type(
2006 &ctx,
2007 "The given view's buffer has been detached and so cannot be used as a response",
2008 ));
2009 }
2010
2011 let objects = ReadableStreamObjects::from_byte_controller(controller).refresh_reader();
2012
2013 ReadableByteStreamController::readable_byte_stream_controller_respond_with_new_view(
2015 ctx, objects, view,
2016 )
2017 }
2018}
2019
2020#[derive(JsLifetime)]
2021pub(super) struct PullIntoDescriptor<'js> {
2022 buffer: ArrayBuffer<'js>,
2023 buffer_byte_length: usize,
2024 byte_offset: usize,
2025 byte_length: usize,
2026 bytes_filled: usize,
2027 minimum_fill: usize,
2028 element_size: usize,
2029 view_constructor: Constructor<'js>,
2030 reader_type: PullIntoDescriptorReaderType,
2031}
2032
2033impl<'js> Trace<'js> for PullIntoDescriptor<'js> {
2034 fn trace<'a>(&self, tracer: Tracer<'a, 'js>) {
2035 self.buffer.trace(tracer);
2036 self.buffer_byte_length.trace(tracer);
2037 self.byte_offset.trace(tracer);
2038 self.byte_length.trace(tracer);
2039 self.bytes_filled.trace(tracer);
2040 self.minimum_fill.trace(tracer);
2041 self.element_size.trace(tracer);
2042 self.view_constructor.trace(tracer);
2043 self.reader_type.trace(tracer);
2044 }
2045}
2046
2047enum PullIntoDescriptorRefMut<'js, 'a> {
2048 Index(usize),
2049 Owned(&'a mut PullIntoDescriptor<'js>),
2050}
2051
2052#[derive(Trace, Clone, Copy)]
2053enum PullIntoDescriptorReaderType {
2054 Default,
2055 Byob,
2056 None,
2057}
2058
2059#[derive(JsLifetime)]
2060struct ReadableByteStreamQueueEntry<'js> {
2061 buffer: ArrayBuffer<'js>,
2062 byte_offset: usize,
2063 byte_length: usize,
2064}
2065
2066impl<'js> Trace<'js> for ReadableByteStreamQueueEntry<'js> {
2067 fn trace<'a>(&self, tracer: Tracer<'a, 'js>) {
2068 self.buffer.trace(tracer);
2069 self.byte_offset.trace(tracer);
2070 self.byte_length.trace(tracer)
2071 }
2072}
2073
2074fn transfer_array_buffer(buffer: ArrayBuffer<'_>) -> Result<ArrayBuffer<'_>> {
2075 buffer.get::<_, Function>("transfer")?.call((This(buffer),))
2076}
2077
2078fn copy_data_block_bytes(
2079 ctx: &Ctx<'_>,
2080 to_block: &ArrayBuffer,
2081 to_index: usize,
2082 from_block: &ArrayBuffer,
2083 from_index: usize,
2084 count: usize,
2085) -> Result<()> {
2086 let to_raw = to_block
2087 .as_raw()
2088 .ok_or(ERROR_MSG_ARRAY_BUFFER_DETACHED)
2089 .or_throw(ctx)?;
2090 let to_slice = unsafe { std::slice::from_raw_parts_mut(to_raw.cast::<u8>().as_ptr(), to_raw.len()) };
2091 let from_raw = from_block
2092 .as_raw()
2093 .ok_or(ERROR_MSG_ARRAY_BUFFER_DETACHED)
2094 .or_throw(ctx)?;
2095 let from_slice = unsafe { std::slice::from_raw_parts(from_raw.cast::<u8>().as_ptr(), from_raw.len()) };
2096
2097 to_slice[to_index..to_index + count]
2098 .copy_from_slice(&from_slice[from_index..from_index + count]);
2099 Ok(())
2100}
2101
2102pub fn readable_byte_stream_controller_enqueue_bytes<'js>(
2106 ctx: Ctx<'js>,
2107 controller: ReadableByteStreamControllerClass<'js>,
2108 buffer: ArrayBuffer<'js>,
2109) -> Result<()> {
2110 readable_byte_stream_controller_enqueue_bytes_inner(ctx, controller, buffer, false)
2111}
2112
2113pub fn readable_byte_stream_controller_enqueue_bytes_borrowed<'js>(
2125 ctx: Ctx<'js>,
2126 controller: ReadableByteStreamControllerClass<'js>,
2127 buffer: ArrayBuffer<'js>,
2128) -> Result<()> {
2129 readable_byte_stream_controller_enqueue_bytes_inner(ctx, controller, buffer, true)
2130}
2131
2132fn readable_byte_stream_controller_enqueue_bytes_inner<'js>(
2133 ctx: Ctx<'js>,
2134 controller: ReadableByteStreamControllerClass<'js>,
2135 buffer: ArrayBuffer<'js>,
2136 skip_transfer: bool,
2137) -> Result<()> {
2138 let byte_length = buffer.len();
2139 if byte_length == 0 {
2140 return Ok(());
2141 }
2142 let view = rquickjs::TypedArray::<u8>::from_arraybuffer(buffer)?;
2143 let borrow = OwnedBorrowMut::from_class(controller);
2144 let chunk = ViewBytes::from_value(
2145 &ctx,
2146 &borrow.function_array_buffer_is_view,
2147 Some(&view.into_value()),
2148 )?;
2149 let objects = ReadableStreamObjects::from_byte_controller(borrow).refresh_reader();
2150 if skip_transfer {
2151 ReadableByteStreamController::readable_byte_stream_controller_enqueue_borrowed(
2152 &ctx, objects, chunk,
2153 )?;
2154 } else {
2155 ReadableByteStreamController::readable_byte_stream_controller_enqueue(
2156 &ctx, objects, chunk,
2157 )?;
2158 }
2159 Ok(())
2160}
2161
2162pub fn readable_byte_stream_controller_close_stream<'js>(
2164 ctx: Ctx<'js>,
2165 controller: ReadableByteStreamControllerClass<'js>,
2166) -> Result<()> {
2167 let borrow = OwnedBorrowMut::from_class(controller);
2168 let objects = ReadableStreamObjects::from_byte_controller(borrow).refresh_reader();
2169 ReadableByteStreamController::readable_byte_stream_controller_close(ctx, objects)?;
2170 Ok(())
2171}