1use std::cell::{Cell, RefCell};
6use std::collections::VecDeque;
7use std::ptr;
8use std::rc::Rc;
9
10use dom_struct::dom_struct;
11use js::context::JSContext;
12use js::jsapi::{Heap, JSObject};
13use js::jsval::{JSVal, UndefinedValue};
14use js::realm::CurrentRealm;
15use js::rust::wrappers2::JS_GetPendingException;
16use js::rust::{HandleObject, HandleValue as SafeHandleValue, HandleValue, MutableHandleValue};
17use js::typedarray::Uint8;
18use script_bindings::conversions::SafeToJSValConvertible;
19use script_bindings::reflector::{Reflector, reflect_dom_object_with_cx};
20
21use crate::dom::bindings::buffer_source::create_buffer_source;
22use crate::dom::bindings::callback::ExceptionHandling;
23use crate::dom::bindings::codegen::Bindings::QueuingStrategyBinding::QueuingStrategySize;
24use crate::dom::bindings::codegen::Bindings::ReadableStreamDefaultControllerBinding::ReadableStreamDefaultControllerMethods;
25use crate::dom::bindings::codegen::UnionTypes::ReadableStreamDefaultControllerOrReadableByteStreamController as Controller;
26use crate::dom::bindings::error::{Error, ErrorToJsval, Fallible, throw_dom_exception};
27use crate::dom::bindings::reflector::DomGlobal;
28use crate::dom::bindings::root::{Dom, DomRoot, MutNullableDom};
29use crate::dom::bindings::trace::RootedTraceableBox;
30use crate::dom::globalscope::GlobalScope;
31use crate::dom::promise::Promise;
32use crate::dom::promisenativehandler::{Callback, PromiseNativeHandler};
33use crate::dom::stream::readablestream::ReadableStream;
34use crate::dom::stream::readablestreamdefaultreader::ReadRequest;
35use crate::dom::stream::underlyingsourcecontainer::{
36 UnderlyingSourceContainer, UnderlyingSourceType,
37};
38use crate::realms::enter_auto_realm;
39
40#[derive(Clone, JSTraceable, MallocSizeOf)]
43#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
44struct PullAlgorithmFulfillmentHandler {
45 controller: Dom<ReadableStreamDefaultController>,
46}
47
48impl Callback for PullAlgorithmFulfillmentHandler {
49 fn callback(&self, cx: &mut CurrentRealm, _v: HandleValue) {
52 self.controller.pulling.set(false);
54
55 if self.controller.pull_again.get() {
57 self.controller.pull_again.set(false);
59
60 self.controller.call_pull_if_needed(cx);
62 }
63 }
64}
65
66#[derive(Clone, JSTraceable, MallocSizeOf)]
69#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
70struct PullAlgorithmRejectionHandler {
71 controller: Dom<ReadableStreamDefaultController>,
72}
73
74impl Callback for PullAlgorithmRejectionHandler {
75 fn callback(&self, cx: &mut CurrentRealm, v: HandleValue) {
78 self.controller.error(cx, v);
80 }
81}
82
83#[derive(Clone, JSTraceable, MallocSizeOf)]
86#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
87struct StartAlgorithmFulfillmentHandler {
88 controller: Dom<ReadableStreamDefaultController>,
89}
90
91impl Callback for StartAlgorithmFulfillmentHandler {
92 fn callback(&self, cx: &mut CurrentRealm, _v: HandleValue) {
95 self.controller.started.set(true);
97
98 self.controller.call_pull_if_needed(cx);
100 }
101}
102
103#[derive(Clone, JSTraceable, MallocSizeOf)]
106#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
107struct StartAlgorithmRejectionHandler {
108 controller: Dom<ReadableStreamDefaultController>,
109}
110
111impl Callback for StartAlgorithmRejectionHandler {
112 fn callback(&self, cx: &mut CurrentRealm, v: HandleValue) {
115 self.controller.error(cx, v);
117 }
118}
119
120#[derive(Debug, JSTraceable, MallocSizeOf, PartialEq)]
122#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
123pub(crate) struct ValueWithSize {
124 #[ignore_malloc_size_of = "Heap is measured by mozjs"]
126 pub(crate) value: Box<Heap<JSVal>>,
127 pub(crate) size: f64,
129}
130
131#[derive(Debug, JSTraceable, MallocSizeOf, PartialEq)]
133#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
134pub(crate) enum EnqueuedValue {
135 Native(Box<[u8]>),
137 Js(ValueWithSize),
139 CloseSentinel,
141}
142
143impl EnqueuedValue {
144 fn size(&self) -> f64 {
145 match self {
146 EnqueuedValue::Native(v) => v.len() as f64,
147 EnqueuedValue::Js(v) => v.size,
148 EnqueuedValue::CloseSentinel => 0.,
151 }
152 }
153
154 fn to_jsval(&self, cx: &mut JSContext, rval: MutableHandleValue) {
155 match self {
156 EnqueuedValue::Native(chunk) => {
157 rooted!(&in(cx) let mut array_buffer_ptr = ptr::null_mut::<JSObject>());
158 create_buffer_source::<Uint8>(cx, chunk, array_buffer_ptr.handle_mut())
159 .expect("failed to create buffer source for native chunk.");
160 array_buffer_ptr.safe_to_jsval(cx, rval);
161 },
162 EnqueuedValue::Js(value_with_size) => value_with_size.value.safe_to_jsval(cx, rval),
163 EnqueuedValue::CloseSentinel => {
164 unreachable!("The close sentinel is never made available as a js val.")
165 },
166 }
167 }
168}
169
170fn is_non_negative_number(value: &EnqueuedValue) -> bool {
172 let value_with_size = match value {
173 EnqueuedValue::Native(_) => return true,
174 EnqueuedValue::Js(value_with_size) => value_with_size,
175 EnqueuedValue::CloseSentinel => return true,
176 };
177
178 if value_with_size.size.is_nan() {
183 return false;
184 }
185
186 if value_with_size.size.is_sign_negative() {
188 return false;
189 }
190
191 true
192}
193
194#[derive(Default, JSTraceable, MallocSizeOf)]
196#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
197pub(crate) struct QueueWithSizes {
198 queue: RefCell<VecDeque<EnqueuedValue>>,
199 pub(crate) total_size: Cell<f64>,
201}
202
203impl QueueWithSizes {
204 pub(crate) fn dequeue_value(&self, cx: &mut JSContext, rval: Option<MutableHandleValue>) {
208 {
209 let queue = self.queue.borrow();
210 let Some(value) = queue.front() else {
211 unreachable!("Buffer cannot be empty when dequeue value is called into.");
212 };
213 self.total_size.set(self.total_size.get() - value.size());
214 if let Some(rval) = rval {
215 value.to_jsval(cx, rval);
216 } else {
217 assert_eq!(value, &EnqueuedValue::CloseSentinel);
218 }
219 }
220 self.queue.borrow_mut().pop_front();
221 }
222
223 #[cfg_attr(crown, expect(crown::unrooted_must_root))]
225 pub(crate) fn enqueue_value_with_size(&self, value: EnqueuedValue) -> Result<(), Error> {
226 if !is_non_negative_number(&value) {
228 return Err(Error::Range(
229 c"The size of the enqueued chunk is not a non-negative number.".to_owned(),
230 ));
231 }
232
233 if value.size().is_infinite() {
235 return Err(Error::Range(
236 c"The size of the enqueued chunk is infinite.".to_owned(),
237 ));
238 }
239
240 self.total_size.set(self.total_size.get() + value.size());
241 self.queue.borrow_mut().push_back(value);
242
243 Ok(())
244 }
245
246 pub(crate) fn is_empty(&self) -> bool {
247 self.queue.borrow().is_empty()
248 }
249
250 pub(crate) fn peek_queue_value(&self, cx: &mut JSContext, rval: MutableHandleValue) -> bool {
253 assert!(!self.is_empty());
258
259 let queue = self.queue.borrow();
261 let value_with_size = queue.front().expect("Queue is not empty.");
262 if let EnqueuedValue::CloseSentinel = value_with_size {
263 return true;
264 }
265
266 value_with_size.to_jsval(cx, rval);
268 false
269 }
270
271 fn get_in_memory_bytes(&self) -> Option<Vec<u8>> {
273 self.queue
274 .borrow()
275 .iter()
276 .try_fold(Vec::new(), |mut acc, value| match value {
277 EnqueuedValue::Native(chunk) => {
278 acc.extend(chunk.iter().copied());
279 Some(acc)
280 },
281 _ => {
282 warn!("get_in_memory_bytes called on a controller with non-native source.");
283 None
284 },
285 })
286 }
287
288 pub(crate) fn reset(&self) {
290 self.queue.borrow_mut().clear();
291 self.total_size.set(Default::default());
292 }
293}
294
295#[dom_struct]
297pub(crate) struct ReadableStreamDefaultController {
298 reflector_: Reflector,
299
300 queue: QueueWithSizes,
302
303 underlying_source: MutNullableDom<UnderlyingSourceContainer>,
309
310 stream: MutNullableDom<ReadableStream>,
311
312 strategy_hwm: f64,
314
315 #[ignore_malloc_size_of = "mozjs"]
317 strategy_size: RefCell<Option<Rc<QueuingStrategySize>>>,
318
319 close_requested: Cell<bool>,
321
322 started: Cell<bool>,
324
325 pulling: Cell<bool>,
327
328 pull_again: Cell<bool>,
330}
331
332impl ReadableStreamDefaultController {
333 fn new_inherited(
334 strategy_hwm: f64,
335 strategy_size: Rc<QueuingStrategySize>,
336 underlying_source: &UnderlyingSourceContainer,
337 ) -> ReadableStreamDefaultController {
338 ReadableStreamDefaultController {
339 reflector_: Reflector::new(),
340 queue: Default::default(),
341 stream: MutNullableDom::new(None),
342 underlying_source: MutNullableDom::new(Some(underlying_source)),
343 strategy_hwm,
344 strategy_size: RefCell::new(Some(strategy_size)),
345 close_requested: Default::default(),
346 started: Default::default(),
347 pulling: Default::default(),
348 pull_again: Default::default(),
349 }
350 }
351
352 pub(crate) fn new(
353 cx: &mut JSContext,
354 global: &GlobalScope,
355 underlying_source: UnderlyingSourceType,
356 strategy_hwm: f64,
357 strategy_size: Rc<QueuingStrategySize>,
358 ) -> DomRoot<ReadableStreamDefaultController> {
359 let underlying_source = UnderlyingSourceContainer::new(cx, global, underlying_source);
360 reflect_dom_object_with_cx(
361 Box::new(ReadableStreamDefaultController::new_inherited(
362 strategy_hwm,
363 strategy_size,
364 &underlying_source,
365 )),
366 global,
367 cx,
368 )
369 }
370
371 pub(crate) fn setup(&self, cx: &mut JSContext, stream: &ReadableStream) -> Result<(), Error> {
373 stream.assert_no_controller();
375
376 self.stream.set(Some(stream));
378
379 let global = &*self.global();
380 let rooted_default_controller = DomRoot::from_ref(self);
381
382 stream.set_default_controller(&rooted_default_controller);
395
396 if let Some(underlying_source) = rooted_default_controller.underlying_source.get() {
397 let start_result = underlying_source
399 .call_start_algorithm(
400 cx,
401 Controller::ReadableStreamDefaultController(rooted_default_controller.clone()),
402 )
403 .unwrap_or_else(|| {
404 let promise = Promise::new_resolved(cx, global, ());
405 Ok(promise)
406 });
407
408 let start_promise = start_result?;
410
411 let handler = PromiseNativeHandler::new(
413 cx,
414 global,
415 Some(Box::new(StartAlgorithmFulfillmentHandler {
416 controller: Dom::from_ref(&rooted_default_controller),
417 })),
418 Some(Box::new(StartAlgorithmRejectionHandler {
419 controller: Dom::from_ref(&rooted_default_controller),
420 })),
421 );
422 let mut realm = enter_auto_realm(cx, global);
423 let cx = &mut realm.current_realm();
424 start_promise.append_native_handler(cx, &handler);
425 };
426
427 Ok(())
428 }
429
430 pub(crate) fn set_underlying_source_this_object(&self, this_object: HandleObject) {
432 if let Some(underlying_source) = self.underlying_source.get() {
433 underlying_source.set_underlying_source_this_object(this_object);
434 }
435 }
436
437 fn dequeue_value(&self, cx: &mut JSContext, rval: MutableHandleValue) {
439 self.queue.dequeue_value(cx, Some(rval));
440 }
441
442 fn should_call_pull(&self) -> bool {
444 let Some(stream) = self.stream.get() else {
448 debug!("`should_call_pull` called on a controller without a stream.");
449 return false;
450 };
451
452 if !self.can_close_or_enqueue() {
454 return false;
455 }
456
457 if !self.started.get() {
459 return false;
460 }
461
462 if stream.is_locked() && stream.get_num_read_requests() > 0 {
465 return true;
466 }
467
468 let desired_size = self.get_desired_size().expect("desiredSize is not null.");
471
472 if desired_size > 0. {
473 return true;
474 }
475
476 false
477 }
478
479 fn call_pull_if_needed(&self, cx: &mut JSContext) {
481 if !self.should_call_pull() {
484 return;
485 }
486
487 if self.pulling.get() {
489 self.pull_again.set(true);
491
492 return;
493 }
494
495 self.pulling.set(true);
497
498 let global = self.global();
501 let rooted_default_controller = DomRoot::from_ref(self);
502 let controller =
503 Controller::ReadableStreamDefaultController(rooted_default_controller.clone());
504
505 let Some(underlying_source) = self.underlying_source.get() else {
506 return;
507 };
508 let handler = PromiseNativeHandler::new(
509 cx,
510 &global,
511 Some(Box::new(PullAlgorithmFulfillmentHandler {
512 controller: Dom::from_ref(&rooted_default_controller),
513 })),
514 Some(Box::new(PullAlgorithmRejectionHandler {
515 controller: Dom::from_ref(&rooted_default_controller),
516 })),
517 );
518
519 let mut realm = enter_auto_realm(cx, &*global);
520 let cx = &mut realm.current_realm();
521
522 let result = underlying_source
523 .call_pull_algorithm(cx, controller)
524 .unwrap_or_else(|| {
525 let promise = Promise::new_resolved(cx, &global, ());
526 Ok(promise)
527 });
528 let promise = result.unwrap_or_else(|error| {
529 rooted!(&in(cx) let mut rval = UndefinedValue());
530 error.to_jsval(cx, &global, rval.handle_mut());
532 Promise::new_rejected(cx, &global, rval.handle())
533 });
534 promise.append_native_handler(cx, &handler);
535 }
536
537 pub(crate) fn perform_cancel_steps(
539 &self,
540 cx: &mut JSContext,
541 global: &GlobalScope,
542 reason: SafeHandleValue,
543 ) -> Rc<Promise> {
544 self.queue.reset();
546
547 let underlying_source = self
548 .underlying_source
549 .get()
550 .expect("Controller should have a source when the cancel steps are called into.");
551 let result = underlying_source
553 .call_cancel_algorithm(cx, global, reason)
554 .unwrap_or_else(|| {
555 let promise = Promise::new(cx, global);
556 promise.resolve_native(cx, &());
557 Ok(promise)
558 });
559 let promise = result.unwrap_or_else(|error| {
560 rooted!(&in(cx) let mut rval = UndefinedValue());
561
562 error.to_jsval(cx, global, rval.handle_mut());
563 let promise = Promise::new(cx, global);
564 promise.reject_native(cx, &rval.handle());
565 promise
566 });
567
568 self.clear_algorithms();
570
571 promise
573 }
574
575 pub(crate) fn perform_pull_steps(&self, cx: &mut JSContext, read_request: &ReadRequest) {
577 let Some(stream) = self.stream.get() else {
580 return;
581 };
582
583 if !self.queue.is_empty() {
585 rooted!(&in(cx) let mut rval = UndefinedValue());
586 let result = RootedTraceableBox::new(Heap::default());
587 self.dequeue_value(cx, rval.handle_mut());
588 result.set(*rval);
589
590 if self.close_requested.get() && self.queue.is_empty() {
592 self.clear_algorithms();
594
595 stream.close(cx);
597 } else {
598 self.call_pull_if_needed(cx);
600 }
601 read_request.chunk_steps(cx, result, &self.global());
603 } else {
604 stream.add_read_request(read_request);
606
607 self.call_pull_if_needed(cx);
609 }
610 }
611
612 pub(crate) fn perform_release_steps(&self) -> Fallible<()> {
614 Ok(())
616 }
617
618 #[expect(unsafe_code)]
620 pub(crate) fn enqueue(&self, cx: &mut JSContext, chunk: SafeHandleValue) -> Result<(), Error> {
621 if !self.can_close_or_enqueue() {
623 return Ok(());
624 }
625
626 let stream = self
627 .stream
628 .get()
629 .expect("Controller must have a stream when a chunk is enqueued.");
630
631 if stream.is_locked() && stream.get_num_read_requests() > 0 {
635 stream.fulfill_read_request(cx, chunk, false);
636 } else {
637 let strategy_size = {
642 let reference = self.strategy_size.borrow();
643 reference.clone()
644 };
645 let size = if let Some(strategy_size) = strategy_size {
646 let result = strategy_size.Call__(cx, chunk, ExceptionHandling::Rethrow);
649 match result {
650 Ok(size) => size,
652 Err(error) => {
653 rooted!(&in(cx) let mut rval = UndefinedValue());
655 unsafe { assert!(JS_GetPendingException(cx, rval.handle_mut())) };
656
657 self.error(cx, rval.handle());
659
660 return Err(error);
663 },
664 }
665 } else {
666 0.
667 };
668
669 {
670 let res = self
672 .queue
673 .enqueue_value_with_size(EnqueuedValue::Js(ValueWithSize {
674 value: Heap::boxed(chunk.get()),
675 size,
676 }));
677 if let Err(error) = res {
678 throw_dom_exception(cx, &self.global(), error);
684
685 rooted!(&in(cx) let mut rval = UndefinedValue());
688 unsafe { assert!(JS_GetPendingException(cx, rval.handle_mut())) };
689
690 self.error(cx, rval.handle());
692
693 return Err(Error::JSFailed);
697 }
698 }
699 }
700
701 self.call_pull_if_needed(cx);
703
704 Ok(())
705 }
706
707 pub(crate) fn enqueue_native(&self, cx: &mut JSContext, chunk: Vec<u8>) {
710 let stream = self
711 .stream
712 .get()
713 .expect("Controller must have a stream when a chunk is enqueued.");
714 if stream.is_locked() && stream.get_num_read_requests() > 0 {
715 rooted!(&in(cx) let mut rval = UndefinedValue());
716 EnqueuedValue::Native(chunk.into_boxed_slice()).to_jsval(cx, rval.handle_mut());
717 stream.fulfill_read_request(cx, rval.handle(), false);
718 } else {
719 self.queue
720 .enqueue_value_with_size(EnqueuedValue::Native(chunk.into_boxed_slice()))
721 .expect("Enqueuing a chunk from Rust should not fail.");
722 }
723 }
724
725 pub(crate) fn in_memory(&self) -> bool {
727 let Some(underlying_source) = self.underlying_source.get() else {
728 return false;
729 };
730 underlying_source.in_memory()
731 }
732
733 pub(crate) fn get_in_memory_bytes(&self) -> Option<Vec<u8>> {
735 let underlying_source = self.underlying_source.get()?;
736 if underlying_source.in_memory() {
737 return self.queue.get_in_memory_bytes();
738 }
739 None
740 }
741
742 fn clear_algorithms(&self) {
744 self.underlying_source.set(None);
747
748 *self.strategy_size.borrow_mut() = None;
750 }
751
752 pub(crate) fn close(&self, cx: &mut JSContext) {
754 if !self.can_close_or_enqueue() {
756 return;
757 }
758
759 let Some(stream) = self.stream.get() else {
760 return;
761 };
762
763 self.close_requested.set(true);
765
766 if self.queue.is_empty() {
767 self.clear_algorithms();
769
770 stream.close(cx);
772 }
773 }
774
775 pub(crate) fn get_desired_size(&self) -> Option<f64> {
777 let stream = self.stream.get()?;
778
779 if stream.is_errored() {
781 return None;
782 }
783
784 if stream.is_closed() {
786 return Some(0.0);
787 }
788
789 let desired_size = self.strategy_hwm - self.queue.total_size.get().clamp(0.0, f64::MAX);
791 Some(desired_size.clamp(desired_size, self.strategy_hwm))
792 }
793
794 pub(crate) fn can_close_or_enqueue(&self) -> bool {
796 let Some(stream) = self.stream.get() else {
797 return false;
798 };
799
800 if !self.close_requested.get() && stream.is_readable() {
802 return true;
803 }
804
805 false
807 }
808
809 pub(crate) fn error(&self, cx: &mut JSContext, e: SafeHandleValue) {
811 let Some(stream) = self.stream.get() else {
812 return;
813 };
814
815 if !stream.is_readable() {
817 return;
818 }
819
820 self.queue.reset();
822
823 self.clear_algorithms();
825
826 stream.error(cx, e);
827 }
828
829 pub(crate) fn has_backpressure(&self) -> bool {
831 !self.should_call_pull()
834 }
835}
836
837impl ReadableStreamDefaultControllerMethods<crate::DomTypeHolder>
838 for ReadableStreamDefaultController
839{
840 fn GetDesiredSize(&self) -> Option<f64> {
842 self.get_desired_size()
843 }
844
845 fn Close(&self, cx: &mut JSContext) -> Fallible<()> {
847 if !self.can_close_or_enqueue() {
848 return Err(Error::Type(c"Stream cannot be closed.".to_owned()));
851 }
852
853 self.close(cx);
855
856 Ok(())
857 }
858
859 fn Enqueue(&self, cx: &mut JSContext, chunk: SafeHandleValue) -> Fallible<()> {
861 if !self.can_close_or_enqueue() {
863 return Err(Error::Type(c"Stream cannot be enqueued to.".to_owned()));
864 }
865
866 self.enqueue(cx, chunk)
868 }
869
870 fn Error(&self, cx: &mut JSContext, e: SafeHandleValue) -> Fallible<()> {
872 self.error(cx, e);
873 Ok(())
874 }
875}