1use std::cell::Cell;
6use std::collections::VecDeque;
7use std::mem;
8use std::rc::Rc;
9
10use dom_struct::dom_struct;
11use js::context::JSContext;
12use js::jsapi::Heap;
13use js::jsval::{JSVal, UndefinedValue};
14use js::realm::CurrentRealm;
15use js::rust::{HandleObject as SafeHandleObject, HandleValue as SafeHandleValue};
16use script_bindings::cell::DomRefCell;
17use script_bindings::reflector::{
18 Reflector, reflect_dom_object_with_cx, reflect_dom_object_with_proto,
19};
20
21use super::byteteereadrequest::ByteTeeReadRequest;
22use super::readablebytestreamcontroller::ReadableByteStreamController;
23use crate::dom::bindings::codegen::Bindings::ReadableStreamDefaultReaderBinding::{
24 ReadableStreamDefaultReaderMethods, ReadableStreamReadResult,
25};
26use crate::dom::bindings::error::{Error, ErrorToJsval, Fallible};
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, TracedPromise};
32use crate::dom::promisenativehandler::{Callback, PromiseNativeHandler};
33use crate::dom::readablestream::{ReadableStream, bytes_from_chunk_jsval};
34use crate::dom::stream::defaultteereadrequest::DefaultTeeReadRequest;
35use crate::dom::stream::readablestreamgenericreader::ReadableStreamGenericReader;
36use crate::dom::types::ReadableStreamDefaultController;
37use crate::realms::enter_auto_realm;
38
39type ReadAllBytesSuccessSteps = dyn Fn(&mut js::context::JSContext, &[u8]);
40type ReadAllBytesFailureSteps = dyn Fn(&mut js::context::JSContext, SafeHandleValue);
41
42impl js::gc::Rootable for ContinueReadMicrotask {}
43
44#[derive(Clone, JSTraceable, MallocSizeOf)]
50#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
51struct ContinueReadMicrotask {
52 reader: Dom<ReadableStreamDefaultReader>,
53 request: ReadRequest,
54}
55
56impl Callback for ContinueReadMicrotask {
57 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
58 self.reader.read(cx, &self.request);
61 }
62}
63
64fn read_loop(
66 cx: &mut js::context::JSContext,
67 reader: &ReadableStreamDefaultReader,
68 success_steps: Rc<ReadAllBytesSuccessSteps>,
69 failure_steps: Rc<ReadAllBytesFailureSteps>,
70) {
71 rooted!(&in(cx) let req = ReadRequest::ReadLoop {
76 success_steps,
77 failure_steps,
78 reader: Dom::from_ref(reader),
79 bytes: Rc::new(DomRefCell::new(Vec::new())),
80 });
81 reader.read(cx, &req);
83}
84
85#[derive(Clone, JSTraceable, MallocSizeOf)]
87#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
88pub(crate) enum ReadRequest {
89 Read(TracedPromise),
91 DefaultTee {
93 tee_read_request: Dom<DefaultTeeReadRequest>,
94 },
95 ReadLoop {
98 #[ignore_malloc_size_of = "dyn Fn"]
99 #[no_trace]
100 success_steps: Rc<ReadAllBytesSuccessSteps>,
101 #[ignore_malloc_size_of = "dyn Fn"]
102 #[no_trace]
103 failure_steps: Rc<ReadAllBytesFailureSteps>,
104 reader: Dom<ReadableStreamDefaultReader>,
105 #[conditional_malloc_size_of]
106 bytes: Rc<DomRefCell<Vec<u8>>>,
107 },
108 ByteTee {
109 byte_tee_read_request: Dom<ByteTeeReadRequest>,
110 },
111}
112
113impl js::rust::Rootable for ReadRequest {}
114
115impl ReadRequest {
116 pub(crate) fn chunk_steps(
118 &self,
119 cx: &mut js::context::JSContext,
120 chunk: RootedTraceableBox<Heap<JSVal>>,
121 global: &GlobalScope,
122 ) {
123 match self {
124 ReadRequest::Read(promise) => {
125 promise.resolve_native(
128 cx,
129 &ReadableStreamReadResult {
130 done: Some(false),
131 value: chunk,
132 },
133 );
134 },
135 ReadRequest::DefaultTee { tee_read_request } => {
136 tee_read_request.enqueue_chunk_steps(cx, chunk);
137 },
138 ReadRequest::ByteTee {
139 byte_tee_read_request,
140 } => {
141 byte_tee_read_request.enqueue_chunk_steps(cx, global, chunk);
142 },
143 ReadRequest::ReadLoop {
144 success_steps: _,
145 failure_steps,
146 reader,
147 bytes,
148 } => {
149 let global = reader.global();
151
152 match bytes_from_chunk_jsval(cx, &chunk) {
153 Ok(vec) => {
154 bytes.borrow_mut().extend_from_slice(&vec);
156
157 let tick = Promise::new(cx, &global);
161 tick.resolve_native(cx, &());
162
163 let handler = PromiseNativeHandler::new(
164 cx,
165 &global,
166 Some(Box::new(ContinueReadMicrotask {
167 reader: Dom::from_ref(reader),
168 request: self.clone(),
169 })),
170 None,
171 );
172
173 let mut realm = enter_auto_realm(cx, &*global);
174 let cx = &mut realm.current_realm();
175 tick.append_native_handler(cx, &handler);
176 },
177 Err(err) => {
178 rooted!(&in(cx) let mut v = UndefinedValue());
180 err.to_jsval(cx, &global, v.handle_mut());
181 (failure_steps)(cx, v.handle());
182 },
183 }
184 },
185 }
186 }
187
188 pub(crate) fn close_steps(&self, cx: &mut js::context::JSContext) {
190 match self {
191 ReadRequest::Read(promise) => {
192 let result = RootedTraceableBox::new(Heap::default());
195 result.set(UndefinedValue());
196 promise.resolve_native(
197 cx,
198 &ReadableStreamReadResult {
199 done: Some(true),
200 value: result,
201 },
202 );
203 },
204 ReadRequest::DefaultTee { tee_read_request } => {
205 tee_read_request.close_steps(cx);
206 },
207 ReadRequest::ByteTee {
208 byte_tee_read_request,
209 } => {
210 byte_tee_read_request
211 .close_steps(cx)
212 .expect("ByteTeeReadRequest close steps should not fail");
213 },
214 ReadRequest::ReadLoop {
215 success_steps,
216 reader,
217 bytes,
218 ..
219 } => {
220 (success_steps)(cx, &bytes.borrow());
222
223 reader
224 .release(cx)
225 .expect("Releasing the read-all-bytes reader should succeed");
226 },
227 }
228 }
229
230 pub(crate) fn error_steps(&self, cx: &mut js::context::JSContext, e: SafeHandleValue) {
232 match self {
233 ReadRequest::Read(promise) => {
234 promise.reject_native(cx, &e)
237 },
238 ReadRequest::DefaultTee { tee_read_request } => {
239 tee_read_request.error_steps();
240 },
241 ReadRequest::ByteTee {
242 byte_tee_read_request,
243 } => {
244 byte_tee_read_request.error_steps();
245 },
246 ReadRequest::ReadLoop {
247 failure_steps,
248 reader,
249 ..
250 } => {
251 (failure_steps)(cx, e);
253
254 reader
255 .release(cx)
256 .expect("Releasing the read-all-bytes reader should succeed");
257 },
258 }
259 }
260}
261
262#[derive(Clone, JSTraceable, MallocSizeOf)]
265#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
266struct ByteTeeClosedPromiseRejectionHandler {
267 branch_1_controller: Dom<ReadableByteStreamController>,
268 branch_2_controller: Dom<ReadableByteStreamController>,
269 #[conditional_malloc_size_of]
270 canceled_1: Rc<Cell<bool>>,
271 #[conditional_malloc_size_of]
272 canceled_2: Rc<Cell<bool>>,
273 #[conditional_malloc_size_of]
274 cancel_promise: Rc<Promise>,
275 #[conditional_malloc_size_of]
276 reader_version: Rc<Cell<u64>>,
277 expected_version: u64,
278}
279
280impl Callback for ByteTeeClosedPromiseRejectionHandler {
281 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
284 if self.reader_version.get() != self.expected_version {
286 return;
287 }
288
289 self.branch_1_controller.error(cx, v);
291
292 self.branch_2_controller.error(cx, v);
294
295 if !self.canceled_1.get() || !self.canceled_2.get() {
297 self.cancel_promise.resolve_native(cx, &());
298 }
299 }
300}
301
302#[derive(Clone, JSTraceable, MallocSizeOf)]
305#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
306struct DefaultTeeClosedPromiseRejectionHandler {
307 branch_1_controller: Dom<ReadableStreamDefaultController>,
308 branch_2_controller: Dom<ReadableStreamDefaultController>,
309 #[conditional_malloc_size_of]
310 canceled_1: Rc<Cell<bool>>,
311 #[conditional_malloc_size_of]
312 canceled_2: Rc<Cell<bool>>,
313 #[conditional_malloc_size_of]
314 cancel_promise: Rc<Promise>,
315}
316
317impl Callback for DefaultTeeClosedPromiseRejectionHandler {
318 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
321 self.branch_1_controller.error(cx, v);
323 self.branch_2_controller.error(cx, v);
325
326 if !self.canceled_1.get() || !self.canceled_2.get() {
328 self.cancel_promise.resolve_native(cx, &());
329 }
330 }
331}
332
333#[dom_struct]
335pub(crate) struct ReadableStreamDefaultReader {
336 reflector_: Reflector,
337
338 stream: MutNullableDom<ReadableStream>,
340
341 read_requests: DomRefCell<VecDeque<ReadRequest>>,
342
343 #[conditional_malloc_size_of]
345 closed_promise: DomRefCell<Rc<Promise>>,
346}
347
348impl ReadableStreamDefaultReader {
349 fn new_with_proto(
350 cx: &mut JSContext,
351 global: &GlobalScope,
352 proto: Option<SafeHandleObject>,
353 ) -> DomRoot<ReadableStreamDefaultReader> {
354 let closed_promise = Promise::new(cx, global);
355 reflect_dom_object_with_proto(
356 cx,
357 Box::new(ReadableStreamDefaultReader::new_inherited(closed_promise)),
358 global,
359 proto,
360 )
361 }
362
363 fn new_inherited(promise: Rc<Promise>) -> ReadableStreamDefaultReader {
364 ReadableStreamDefaultReader {
365 reflector_: Reflector::new(),
366 stream: MutNullableDom::new(None),
367 read_requests: DomRefCell::new(Default::default()),
368 closed_promise: DomRefCell::new(promise),
369 }
370 }
371
372 pub(crate) fn new(
373 cx: &mut JSContext,
374 global: &GlobalScope,
375 ) -> DomRoot<ReadableStreamDefaultReader> {
376 let closed_promise = Promise::new(cx, global);
377 reflect_dom_object_with_cx(Box::new(Self::new_inherited(closed_promise)), global, cx)
378 }
379
380 pub(crate) fn set_up(
382 &self,
383 cx: &mut JSContext,
384 stream: &ReadableStream,
385 global: &GlobalScope,
386 ) -> Fallible<()> {
387 if stream.is_locked() {
389 return Err(Error::Type(c"stream is locked".to_owned()));
390 }
391 self.generic_initialize(cx, global, stream);
394
395 self.read_requests.borrow_mut().clear();
397
398 Ok(())
399 }
400
401 pub(crate) fn close(&self, cx: &mut js::context::JSContext) {
403 self.closed_promise.borrow().resolve_native(cx, &());
405 rooted!(&in(cx) let mut read_requests = self.take_read_requests());
408 for request in read_requests.iter() {
411 request.close_steps(cx);
413 }
414 }
415
416 pub(crate) fn add_read_request(&self, read_request: &ReadRequest) {
418 self.read_requests
419 .borrow_mut()
420 .push_back(read_request.clone());
421 }
422
423 pub(crate) fn get_num_read_requests(&self) -> usize {
425 self.read_requests.borrow().len()
426 }
427
428 pub(crate) fn error(&self, cx: &mut js::context::JSContext, e: SafeHandleValue) {
430 self.closed_promise.borrow().reject_native(cx, &e);
432
433 self.closed_promise.borrow().set_promise_is_handled(cx);
435
436 self.error_read_requests(cx, e);
438 }
439
440 pub(crate) fn remove_read_request(&self) -> ReadRequest {
442 self.read_requests
443 .borrow_mut()
444 .pop_front()
445 .expect("Reader must have read request when remove is called into.")
446 }
447
448 pub(crate) fn release(&self, cx: &mut js::context::JSContext) -> Fallible<()> {
450 self.generic_release(cx).expect("Generic release failed");
452 rooted!(&in(cx) let mut error = UndefinedValue());
454 Error::Type(c"Reader is released".to_owned()).to_jsval(
455 cx,
456 &self.global(),
457 error.handle_mut(),
458 );
459
460 self.error_read_requests(cx, error.handle());
462 Ok(())
463 }
464
465 fn take_read_requests(&self) -> VecDeque<ReadRequest> {
466 mem::take(&mut *self.read_requests.borrow_mut())
467 }
468
469 fn error_read_requests(&self, cx: &mut js::context::JSContext, rval: SafeHandleValue) {
471 rooted!(&in(cx) let mut read_requests = self.take_read_requests());
473
474 for request in read_requests.iter() {
476 request.error_steps(cx, rval);
477 }
478 }
479
480 pub(crate) fn read(&self, cx: &mut js::context::JSContext, read_request: &ReadRequest) {
482 assert!(self.stream.get().is_some());
486
487 let stream = self.stream.get().unwrap();
488
489 stream.set_is_disturbed(true);
491 if stream.is_closed() {
493 read_request.close_steps(cx);
494 } else if stream.is_errored() {
495 rooted!(&in(cx) let mut error = UndefinedValue());
498 stream.get_stored_error(error.handle_mut());
499 read_request.error_steps(cx, error.handle());
500 } else {
501 assert!(stream.is_readable());
504 stream.perform_pull_steps(cx, read_request);
506 }
507 }
508
509 #[allow(clippy::too_many_arguments)]
512 pub(crate) fn byte_tee_append_native_handler_to_closed_promise(
513 &self,
514 cx: &mut js::context::JSContext,
515 branch_1: &ReadableStream,
516 branch_2: &ReadableStream,
517 canceled_1: Rc<Cell<bool>>,
518 canceled_2: Rc<Cell<bool>>,
519 cancel_promise: Rc<Promise>,
520 reader_version: Rc<Cell<u64>>,
521 expected_version: u64,
522 ) {
523 let branch_1_controller = branch_1.get_byte_controller();
525 let branch_2_controller = branch_2.get_byte_controller();
526
527 let global = self.global();
528 let handler = PromiseNativeHandler::new(
529 cx,
530 &global,
531 None,
532 Some(Box::new(ByteTeeClosedPromiseRejectionHandler {
533 branch_1_controller: Dom::from_ref(&branch_1_controller),
534 branch_2_controller: Dom::from_ref(&branch_2_controller),
535 canceled_1,
536 canceled_2,
537 cancel_promise,
538 reader_version,
539 expected_version,
540 })),
541 );
542
543 let mut realm = enter_auto_realm(cx, &*global);
544 let cx = &mut realm.current_realm();
545
546 self.closed_promise
547 .borrow()
548 .append_native_handler(cx, &handler);
549 }
550
551 pub(crate) fn default_tee_append_native_handler_to_closed_promise(
553 &self,
554 cx: &mut js::context::JSContext,
555 branch_1: &ReadableStream,
556 branch_2: &ReadableStream,
557 canceled_1: Rc<Cell<bool>>,
558 canceled_2: Rc<Cell<bool>>,
559 cancel_promise: Rc<Promise>,
560 ) {
561 let branch_1_controller = branch_1.get_default_controller();
562
563 let branch_2_controller = branch_2.get_default_controller();
564
565 let global = self.global();
566 let handler = PromiseNativeHandler::new(
567 cx,
568 &global,
569 None,
570 Some(Box::new(DefaultTeeClosedPromiseRejectionHandler {
571 branch_1_controller: Dom::from_ref(&branch_1_controller),
572 branch_2_controller: Dom::from_ref(&branch_2_controller),
573 canceled_1,
574 canceled_2,
575 cancel_promise,
576 })),
577 );
578
579 let mut realm = enter_auto_realm(cx, &*global);
580 let cx = &mut realm.current_realm();
581
582 self.closed_promise
583 .borrow()
584 .append_native_handler(cx, &handler);
585 }
586
587 pub(crate) fn read_all_bytes(
589 &self,
590 cx: &mut js::context::JSContext,
591 success_steps: Rc<ReadAllBytesSuccessSteps>,
592 failure_steps: Rc<ReadAllBytesFailureSteps>,
593 ) {
594 read_loop(cx, self, success_steps, failure_steps);
599 }
600
601 pub(crate) fn process_read_requests(
603 &self,
604 cx: &mut js::context::JSContext,
605 controller: &ReadableByteStreamController,
606 ) -> Fallible<()> {
607 while !self.read_requests.borrow().is_empty() {
609 if controller.get_queue_total_size() == 0.0 {
611 return Ok(());
612 }
613
614 rooted!(&in(cx) let read_request = self.remove_read_request());
617
618 controller
620 .fill_read_request_from_queue(cx, &read_request)
621 .expect("Fill read request from queue failed");
622 }
623 Ok(())
624 }
625}
626
627impl ReadableStreamDefaultReaderMethods<crate::DomTypeHolder> for ReadableStreamDefaultReader {
628 fn Constructor(
630 cx: &mut JSContext,
631 global: &GlobalScope,
632 proto: Option<SafeHandleObject>,
633 stream: &ReadableStream,
634 ) -> Fallible<DomRoot<Self>> {
635 let reader = Self::new_with_proto(cx, global, proto);
636
637 reader.set_up(cx, stream, global)?;
639
640 Ok(reader)
641 }
642
643 fn Read(&self, cx: &mut js::context::JSContext) -> Rc<Promise> {
645 if self.stream.get().is_none() {
647 rooted!(&in(cx) let mut error = UndefinedValue());
648 Error::Type(c"stream is undefined".to_owned()).to_jsval(
649 cx,
650 &self.global(),
651 error.handle_mut(),
652 );
653 return Promise::new_rejected(cx, &self.global(), error.handle());
654 }
655 let promise = Promise::new_rooted(cx, &self.global());
657
658 rooted!(&in(cx) let read_request = ReadRequest::Read(promise.to_traced()));
669
670 self.read(cx, &read_request);
672
673 promise.into()
675 }
676
677 fn ReleaseLock(&self, cx: &mut js::context::JSContext) -> Fallible<()> {
679 if self.stream.get().is_none() {
680 return Ok(());
682 }
683
684 self.release(cx)
686 }
687
688 fn Closed(&self) -> Rc<Promise> {
690 self.closed()
691 }
692
693 fn Cancel(&self, cx: &mut js::context::JSContext, reason: SafeHandleValue) -> Rc<Promise> {
695 self.generic_cancel(cx, &self.global(), reason)
696 }
697}
698
699impl ReadableStreamGenericReader for ReadableStreamDefaultReader {
700 fn get_closed_promise(&self) -> Rc<Promise> {
701 self.closed_promise.borrow().clone()
702 }
703
704 fn set_closed_promise(&self, promise: Rc<Promise>) {
705 *self.closed_promise.borrow_mut() = promise;
706 }
707
708 fn set_stream(&self, stream: Option<&ReadableStream>) {
709 self.stream.set(stream);
710 }
711
712 fn get_stream(&self) -> Option<DomRoot<ReadableStream>> {
713 self.stream.get()
714 }
715
716 fn as_default_reader(&self) -> Option<&ReadableStreamDefaultReader> {
717 Some(self)
718 }
719}