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, RootedPromise, 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 cancel_promise: TracedPromise,
274 #[conditional_malloc_size_of]
275 reader_version: Rc<Cell<u64>>,
276 expected_version: u64,
277}
278
279impl Callback for ByteTeeClosedPromiseRejectionHandler {
280 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
283 if self.reader_version.get() != self.expected_version {
285 return;
286 }
287
288 self.branch_1_controller.error(cx, v);
290
291 self.branch_2_controller.error(cx, v);
293
294 if !self.canceled_1.get() || !self.canceled_2.get() {
296 self.cancel_promise.resolve_native(cx, &());
297 }
298 }
299}
300
301#[derive(Clone, JSTraceable, MallocSizeOf)]
304#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
305struct DefaultTeeClosedPromiseRejectionHandler {
306 branch_1_controller: Dom<ReadableStreamDefaultController>,
307 branch_2_controller: Dom<ReadableStreamDefaultController>,
308 #[conditional_malloc_size_of]
309 canceled_1: Rc<Cell<bool>>,
310 #[conditional_malloc_size_of]
311 canceled_2: Rc<Cell<bool>>,
312 cancel_promise: TracedPromise,
313}
314
315impl Callback for DefaultTeeClosedPromiseRejectionHandler {
316 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
319 self.branch_1_controller.error(cx, v);
321 self.branch_2_controller.error(cx, v);
323
324 if !self.canceled_1.get() || !self.canceled_2.get() {
326 self.cancel_promise.resolve_native(cx, &());
327 }
328 }
329}
330
331#[dom_struct]
333pub(crate) struct ReadableStreamDefaultReader {
334 reflector_: Reflector,
335
336 stream: MutNullableDom<ReadableStream>,
338
339 read_requests: DomRefCell<VecDeque<ReadRequest>>,
340
341 closed_promise: DomRefCell<TracedPromise>,
343}
344
345impl ReadableStreamDefaultReader {
346 fn new_with_proto(
347 cx: &mut JSContext,
348 global: &GlobalScope,
349 proto: Option<SafeHandleObject>,
350 ) -> DomRoot<ReadableStreamDefaultReader> {
351 let closed_promise = Promise::new_rooted(cx, global);
352 reflect_dom_object_with_proto(
353 cx,
354 Box::new(ReadableStreamDefaultReader::new_inherited(&closed_promise)),
355 global,
356 proto,
357 )
358 }
359
360 fn new_inherited(promise: &RootedPromise) -> ReadableStreamDefaultReader {
361 ReadableStreamDefaultReader {
362 reflector_: Reflector::new(),
363 stream: MutNullableDom::new(None),
364 read_requests: DomRefCell::new(Default::default()),
365 closed_promise: DomRefCell::new(promise.to_traced()),
366 }
367 }
368
369 pub(crate) fn new(
370 cx: &mut JSContext,
371 global: &GlobalScope,
372 ) -> DomRoot<ReadableStreamDefaultReader> {
373 let closed_promise = Promise::new_rooted(cx, global);
374 reflect_dom_object_with_cx(Box::new(Self::new_inherited(&closed_promise)), global, cx)
375 }
376
377 pub(crate) fn set_up(
379 &self,
380 cx: &mut JSContext,
381 stream: &ReadableStream,
382 global: &GlobalScope,
383 ) -> Fallible<()> {
384 if stream.is_locked() {
386 return Err(Error::Type(c"stream is locked".to_owned()));
387 }
388 self.generic_initialize(cx, global, stream);
391
392 self.read_requests.borrow_mut().clear();
394
395 Ok(())
396 }
397
398 pub(crate) fn close(&self, cx: &mut js::context::JSContext) {
400 self.closed_promise.borrow().resolve_native(cx, &());
402 rooted!(&in(cx) let mut read_requests = self.take_read_requests());
405 for request in read_requests.iter() {
408 request.close_steps(cx);
410 }
411 }
412
413 pub(crate) fn add_read_request(&self, read_request: &ReadRequest) {
415 self.read_requests
416 .borrow_mut()
417 .push_back(read_request.clone());
418 }
419
420 pub(crate) fn get_num_read_requests(&self) -> usize {
422 self.read_requests.borrow().len()
423 }
424
425 pub(crate) fn error(&self, cx: &mut js::context::JSContext, e: SafeHandleValue) {
427 self.closed_promise.borrow().reject_native(cx, &e);
429
430 self.closed_promise.borrow().set_promise_is_handled(cx);
432
433 self.error_read_requests(cx, e);
435 }
436
437 pub(crate) fn remove_read_request(&self) -> ReadRequest {
439 self.read_requests
440 .borrow_mut()
441 .pop_front()
442 .expect("Reader must have read request when remove is called into.")
443 }
444
445 pub(crate) fn release(&self, cx: &mut js::context::JSContext) -> Fallible<()> {
447 self.generic_release(cx).expect("Generic release failed");
449 rooted!(&in(cx) let mut error = UndefinedValue());
451 Error::Type(c"Reader is released".to_owned()).to_jsval(
452 cx,
453 &self.global(),
454 error.handle_mut(),
455 );
456
457 self.error_read_requests(cx, error.handle());
459 Ok(())
460 }
461
462 fn take_read_requests(&self) -> VecDeque<ReadRequest> {
463 mem::take(&mut *self.read_requests.borrow_mut())
464 }
465
466 fn error_read_requests(&self, cx: &mut js::context::JSContext, rval: SafeHandleValue) {
468 rooted!(&in(cx) let mut read_requests = self.take_read_requests());
470
471 for request in read_requests.iter() {
473 request.error_steps(cx, rval);
474 }
475 }
476
477 pub(crate) fn read(&self, cx: &mut js::context::JSContext, read_request: &ReadRequest) {
479 assert!(self.stream.get().is_some());
483
484 let stream = self.stream.get().unwrap();
485
486 stream.set_is_disturbed(true);
488 if stream.is_closed() {
490 read_request.close_steps(cx);
491 } else if stream.is_errored() {
492 rooted!(&in(cx) let mut error = UndefinedValue());
495 stream.get_stored_error(error.handle_mut());
496 read_request.error_steps(cx, error.handle());
497 } else {
498 assert!(stream.is_readable());
501 stream.perform_pull_steps(cx, read_request);
503 }
504 }
505
506 #[allow(clippy::too_many_arguments)]
509 pub(crate) fn byte_tee_append_native_handler_to_closed_promise(
510 &self,
511 cx: &mut js::context::JSContext,
512 branch_1: &ReadableStream,
513 branch_2: &ReadableStream,
514 canceled_1: Rc<Cell<bool>>,
515 canceled_2: Rc<Cell<bool>>,
516 cancel_promise: &RootedPromise,
517 reader_version: Rc<Cell<u64>>,
518 expected_version: u64,
519 ) {
520 let branch_1_controller = branch_1.get_byte_controller();
522 let branch_2_controller = branch_2.get_byte_controller();
523
524 let global = self.global();
525 let handler = PromiseNativeHandler::new(
526 cx,
527 &global,
528 None,
529 Some(Box::new(ByteTeeClosedPromiseRejectionHandler {
530 branch_1_controller: Dom::from_ref(&branch_1_controller),
531 branch_2_controller: Dom::from_ref(&branch_2_controller),
532 canceled_1,
533 canceled_2,
534 cancel_promise: cancel_promise.to_traced(),
535 reader_version,
536 expected_version,
537 })),
538 );
539
540 let mut realm = enter_auto_realm(cx, &*global);
541 let cx = &mut realm.current_realm();
542
543 self.closed_promise
544 .borrow()
545 .append_native_handler(cx, &handler);
546 }
547
548 pub(crate) fn default_tee_append_native_handler_to_closed_promise(
550 &self,
551 cx: &mut js::context::JSContext,
552 branch_1: &ReadableStream,
553 branch_2: &ReadableStream,
554 canceled_1: Rc<Cell<bool>>,
555 canceled_2: Rc<Cell<bool>>,
556 cancel_promise: &RootedPromise,
557 ) {
558 let branch_1_controller = branch_1.get_default_controller();
559
560 let branch_2_controller = branch_2.get_default_controller();
561
562 let global = self.global();
563 let handler = PromiseNativeHandler::new(
564 cx,
565 &global,
566 None,
567 Some(Box::new(DefaultTeeClosedPromiseRejectionHandler {
568 branch_1_controller: Dom::from_ref(&branch_1_controller),
569 branch_2_controller: Dom::from_ref(&branch_2_controller),
570 canceled_1,
571 canceled_2,
572 cancel_promise: cancel_promise.to_traced(),
573 })),
574 );
575
576 let mut realm = enter_auto_realm(cx, &*global);
577 let cx = &mut realm.current_realm();
578
579 self.closed_promise
580 .borrow()
581 .append_native_handler(cx, &handler);
582 }
583
584 pub(crate) fn read_all_bytes(
586 &self,
587 cx: &mut js::context::JSContext,
588 success_steps: Rc<ReadAllBytesSuccessSteps>,
589 failure_steps: Rc<ReadAllBytesFailureSteps>,
590 ) {
591 read_loop(cx, self, success_steps, failure_steps);
596 }
597
598 pub(crate) fn process_read_requests(
600 &self,
601 cx: &mut js::context::JSContext,
602 controller: &ReadableByteStreamController,
603 ) -> Fallible<()> {
604 while !self.read_requests.borrow().is_empty() {
606 if controller.get_queue_total_size() == 0.0 {
608 return Ok(());
609 }
610
611 rooted!(&in(cx) let read_request = self.remove_read_request());
614
615 controller
617 .fill_read_request_from_queue(cx, &read_request)
618 .expect("Fill read request from queue failed");
619 }
620 Ok(())
621 }
622}
623
624impl ReadableStreamDefaultReaderMethods<crate::DomTypeHolder> for ReadableStreamDefaultReader {
625 fn Constructor(
627 cx: &mut JSContext,
628 global: &GlobalScope,
629 proto: Option<SafeHandleObject>,
630 stream: &ReadableStream,
631 ) -> Fallible<DomRoot<Self>> {
632 let reader = Self::new_with_proto(cx, global, proto);
633
634 reader.set_up(cx, stream, global)?;
636
637 Ok(reader)
638 }
639
640 fn Read(&self, cx: &mut js::context::JSContext) -> RootedPromise {
642 if self.stream.get().is_none() {
644 rooted!(&in(cx) let mut error = UndefinedValue());
645 Error::Type(c"stream is undefined".to_owned()).to_jsval(
646 cx,
647 &self.global(),
648 error.handle_mut(),
649 );
650 return Promise::new_rejected_rooted(cx, &self.global(), error.handle());
651 }
652 let promise = Promise::new_rooted(cx, &self.global());
654
655 rooted!(&in(cx) let read_request = ReadRequest::Read(promise.to_traced()));
666
667 self.read(cx, &read_request);
669
670 promise
672 }
673
674 fn ReleaseLock(&self, cx: &mut js::context::JSContext) -> Fallible<()> {
676 if self.stream.get().is_none() {
677 return Ok(());
679 }
680
681 self.release(cx)
683 }
684
685 fn Closed(&self, cx: &JSContext) -> RootedPromise {
687 self.closed(cx)
688 }
689
690 fn Cancel(&self, cx: &mut js::context::JSContext, reason: SafeHandleValue) -> RootedPromise {
692 self.generic_cancel(cx, &self.global(), reason)
693 }
694}
695
696impl ReadableStreamGenericReader for ReadableStreamDefaultReader {
697 fn get_closed_promise(&self, cx: &JSContext) -> RootedPromise {
698 self.closed_promise.borrow().root(cx)
699 }
700
701 fn set_closed_promise(&self, promise: &RootedPromise) {
702 *self.closed_promise.borrow_mut() = promise.to_traced();
703 }
704
705 fn set_stream(&self, stream: Option<&ReadableStream>) {
706 self.stream.set(stream);
707 }
708
709 fn get_stream(&self) -> Option<DomRoot<ReadableStream>> {
710 self.stream.get()
711 }
712
713 fn as_default_reader(&self) -> Option<&ReadableStreamDefaultReader> {
714 Some(self)
715 }
716}