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::gc::CustomAutoRooterGuard;
13use js::jsapi::Heap;
14use js::jsval::{JSVal, UndefinedValue};
15use js::realm::CurrentRealm;
16use js::rust::{HandleObject as SafeHandleObject, HandleValue as SafeHandleValue};
17use js::typedarray::{ArrayBufferView, ArrayBufferViewU8};
18use script_bindings::cell::DomRefCell;
19use script_bindings::reflector::{
20 Reflector, reflect_dom_object_with_cx, reflect_dom_object_with_proto,
21};
22use script_bindings::root::Dom;
23
24use super::byteteereadintorequest::ByteTeeReadIntoRequest;
25use super::readablebytestreamcontroller::ReadableByteStreamController;
26use super::readablestreamgenericreader::ReadableStreamGenericReader;
27use crate::dom::bindings::buffer_source::HeapBufferSource;
28use crate::dom::bindings::codegen::Bindings::ReadableStreamBYOBReaderBinding::{
29 ReadableStreamBYOBReaderMethods, ReadableStreamBYOBReaderReadOptions,
30};
31use crate::dom::bindings::codegen::Bindings::ReadableStreamDefaultReaderBinding::ReadableStreamReadResult;
32use crate::dom::bindings::error::{Error, ErrorToJsval, Fallible};
33use crate::dom::bindings::reflector::DomGlobal;
34use crate::dom::bindings::root::{DomRoot, MutNullableDom};
35use crate::dom::bindings::trace::RootedTraceableBox;
36use crate::dom::globalscope::GlobalScope;
37use crate::dom::promise::Promise;
38use crate::dom::promisenativehandler::{Callback, PromiseNativeHandler};
39use crate::dom::stream::readablestream::ReadableStream;
40use crate::realms::enter_auto_realm;
41
42#[derive(Clone, JSTraceable, MallocSizeOf)]
44pub enum ReadIntoRequest {
45 Read(#[conditional_malloc_size_of] Rc<Promise>),
47 ByteTee {
48 byte_tee_read_into_request: Dom<ByteTeeReadIntoRequest>,
49 },
50}
51
52impl ReadIntoRequest {
53 pub fn chunk_steps(&self, cx: &mut JSContext, chunk: RootedTraceableBox<Heap<JSVal>>) {
55 match self {
56 ReadIntoRequest::Read(promise) => {
57 promise.resolve_native(
60 cx,
61 &ReadableStreamReadResult {
62 done: Some(false),
63 value: chunk,
64 },
65 );
66 },
67 ReadIntoRequest::ByteTee {
68 byte_tee_read_into_request,
69 } => {
70 rooted!(&in(cx) let chunk_object = chunk.get().to_object());
71 byte_tee_read_into_request.enqueue_chunk_steps(
72 cx,
73 RootedTraceableBox::new(HeapBufferSource::<ArrayBufferViewU8>::new(
74 chunk_object.handle(),
75 )),
76 )
77 },
78 }
79 }
80
81 pub fn close_steps(&self, cx: &mut JSContext, chunk: Option<RootedTraceableBox<Heap<JSVal>>>) {
83 match self {
84 ReadIntoRequest::Read(promise) => match chunk {
85 Some(chunk) => promise.resolve_native(
88 cx,
89 &ReadableStreamReadResult {
90 done: Some(true),
91 value: chunk,
92 },
93 ),
94 None => {
95 let result = RootedTraceableBox::new(Heap::default());
96 result.set(UndefinedValue());
97 promise.resolve_native(
98 cx,
99 &ReadableStreamReadResult {
100 done: Some(true),
101 value: result,
102 },
103 );
104 },
105 },
106 ReadIntoRequest::ByteTee {
107 byte_tee_read_into_request,
108 } => match chunk {
109 Some(chunk) => {
110 rooted!(&in(cx) let chunk_object = chunk.get().to_object());
111 byte_tee_read_into_request
112 .close_steps(
113 cx,
114 Some(RootedTraceableBox::new(
115 HeapBufferSource::<ArrayBufferViewU8>::new(chunk_object.handle()),
116 )),
117 )
118 .expect("close steps should not fail")
119 },
120 None => byte_tee_read_into_request
121 .close_steps(cx, None)
122 .expect("close steps should not fail"),
123 },
124 }
125 }
126
127 pub(crate) fn error_steps(&self, cx: &mut JSContext, e: SafeHandleValue) {
129 match self {
130 ReadIntoRequest::Read(promise) => {
131 promise.reject_native(cx, &e)
134 },
135 ReadIntoRequest::ByteTee {
136 byte_tee_read_into_request,
137 } => {
138 byte_tee_read_into_request.error_steps();
139 },
140 }
141 }
142}
143
144#[derive(Clone, JSTraceable, MallocSizeOf)]
147#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
148struct ByteTeeClosedPromiseRejectionHandler {
149 branch_1_controller: Dom<ReadableByteStreamController>,
150 branch_2_controller: Dom<ReadableByteStreamController>,
151 #[conditional_malloc_size_of]
152 canceled_1: Rc<Cell<bool>>,
153 #[conditional_malloc_size_of]
154 canceled_2: Rc<Cell<bool>>,
155 #[conditional_malloc_size_of]
156 cancel_promise: Rc<Promise>,
157 #[conditional_malloc_size_of]
158 reader_version: Rc<Cell<u64>>,
159 expected_version: u64,
160}
161
162impl Callback for ByteTeeClosedPromiseRejectionHandler {
163 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
166 if self.reader_version.get() != self.expected_version {
168 return;
169 }
170
171 self.branch_1_controller.error(cx, v);
173
174 self.branch_2_controller.error(cx, v);
176
177 if !self.canceled_1.get() || !self.canceled_2.get() {
179 self.cancel_promise.resolve_native(cx, &());
180 }
181 }
182}
183
184#[dom_struct]
186pub(crate) struct ReadableStreamBYOBReader {
187 reflector_: Reflector,
188
189 stream: MutNullableDom<ReadableStream>,
191
192 read_into_requests: DomRefCell<VecDeque<ReadIntoRequest>>,
193
194 #[conditional_malloc_size_of]
196 closed_promise: DomRefCell<Rc<Promise>>,
197}
198
199impl ReadableStreamBYOBReader {
200 fn new_with_proto(
201 cx: &mut JSContext,
202 global: &GlobalScope,
203 proto: Option<SafeHandleObject>,
204 ) -> DomRoot<ReadableStreamBYOBReader> {
205 let closed_promise = Promise::new(cx, global);
206 reflect_dom_object_with_proto(
207 cx,
208 Box::new(ReadableStreamBYOBReader::new_inherited(closed_promise)),
209 global,
210 proto,
211 )
212 }
213
214 fn new_inherited(promise: Rc<Promise>) -> ReadableStreamBYOBReader {
215 ReadableStreamBYOBReader {
216 reflector_: Reflector::new(),
217 stream: MutNullableDom::new(None),
218 read_into_requests: DomRefCell::new(Default::default()),
219 closed_promise: DomRefCell::new(promise),
220 }
221 }
222
223 pub(crate) fn new(
224 cx: &mut JSContext,
225 global: &GlobalScope,
226 ) -> DomRoot<ReadableStreamBYOBReader> {
227 let closed_promise = Promise::new(cx, global);
228 reflect_dom_object_with_cx(Box::new(Self::new_inherited(closed_promise)), global, cx)
229 }
230
231 pub(crate) fn set_up(
233 &self,
234 cx: &mut JSContext,
235 stream: &ReadableStream,
236 global: &GlobalScope,
237 ) -> Fallible<()> {
238 if stream.is_locked() {
240 return Err(Error::Type(c"stream is locked".to_owned()));
241 }
242
243 if !stream.has_byte_controller() {
245 return Err(Error::Type(
246 c"stream controller is not a byte stream controller".to_owned(),
247 ));
248 }
249
250 self.generic_initialize(cx, global, stream);
252
253 self.read_into_requests.borrow_mut().clear();
255
256 Ok(())
257 }
258
259 pub(crate) fn release(&self, cx: &mut JSContext) -> Fallible<()> {
261 self.generic_release(cx).expect("Generic release failed");
263 rooted!(&in(cx) let mut error = UndefinedValue());
265 Error::Type(c"Reader is released".to_owned()).to_jsval(
266 cx,
267 &self.global(),
268 error.handle_mut(),
269 );
270
271 self.error_read_into_requests(cx, error.handle());
273 Ok(())
274 }
275
276 pub(crate) fn error_read_into_requests(&self, cx: &mut JSContext, e: SafeHandleValue) {
278 self.closed_promise.borrow().reject_native(cx, &e);
280
281 self.closed_promise.borrow().set_promise_is_handled(cx);
283
284 let mut read_into_requests = self.take_read_into_requests();
286
287 for request in read_into_requests.drain(0..) {
289 request.error_steps(cx, e);
291 }
292 }
293
294 fn take_read_into_requests(&self) -> VecDeque<ReadIntoRequest> {
295 mem::take(&mut *self.read_into_requests.borrow_mut())
296 }
297
298 pub(crate) fn add_read_into_request(&self, read_request: &ReadIntoRequest) {
300 self.read_into_requests
301 .borrow_mut()
302 .push_back(read_request.clone());
303 }
304
305 pub(crate) fn cancel(&self, cx: &mut JSContext) {
307 let mut read_into_requests = self.take_read_into_requests();
310 for request in read_into_requests.drain(0..) {
313 request.close_steps(cx, None);
315 }
316 }
317
318 pub(crate) fn close(&self, cx: &mut JSContext) {
319 self.closed_promise.borrow().resolve_native(cx, &());
321 }
322
323 pub(crate) fn read(
325 &self,
326 cx: &mut JSContext,
327 view: &HeapBufferSource<ArrayBufferViewU8>,
328 min: u64,
329 read_into_request: &ReadIntoRequest,
330 ) {
331 assert!(self.stream.get().is_some());
335
336 let stream = self.stream.get().unwrap();
337
338 stream.set_is_disturbed(true);
340 if stream.is_errored() {
342 rooted!(&in(cx) let mut error = UndefinedValue());
343 stream.get_stored_error(error.handle_mut());
344
345 read_into_request.error_steps(cx, error.handle());
346 } else {
347 stream.perform_pull_into(cx, read_into_request, view, min);
350 }
351 }
352
353 pub(crate) fn get_num_read_into_requests(&self) -> usize {
354 self.read_into_requests.borrow().len()
355 }
356
357 pub(crate) fn remove_read_into_request(&self) -> ReadIntoRequest {
358 self.read_into_requests
359 .borrow_mut()
360 .pop_front()
361 .expect("read into requests is empty")
362 }
363
364 #[allow(clippy::too_many_arguments)]
365 pub(crate) fn byte_tee_append_native_handler_to_closed_promise(
366 &self,
367 cx: &mut JSContext,
368 branch_1: &ReadableStream,
369 branch_2: &ReadableStream,
370 canceled_1: Rc<Cell<bool>>,
371 canceled_2: Rc<Cell<bool>>,
372 cancel_promise: Rc<Promise>,
373 reader_version: Rc<Cell<u64>>,
374 expected_version: u64,
375 ) {
376 let branch_1_controller = branch_1.get_byte_controller();
377
378 let branch_2_controller = branch_2.get_byte_controller();
379
380 let global = self.global();
381 let handler = PromiseNativeHandler::new(
382 cx,
383 &global,
384 None,
385 Some(Box::new(ByteTeeClosedPromiseRejectionHandler {
386 branch_1_controller: Dom::from_ref(&branch_1_controller),
387 branch_2_controller: Dom::from_ref(&branch_2_controller),
388 canceled_1,
389 canceled_2,
390 cancel_promise,
391 reader_version,
392 expected_version,
393 })),
394 );
395
396 let mut realm = enter_auto_realm(cx, &*global);
397 let cx = &mut realm.current_realm();
398
399 self.closed_promise
400 .borrow()
401 .append_native_handler(cx, &handler);
402 }
403}
404
405impl ReadableStreamBYOBReaderMethods<crate::DomTypeHolder> for ReadableStreamBYOBReader {
406 fn Constructor(
408 cx: &mut JSContext,
409 global: &GlobalScope,
410 proto: Option<SafeHandleObject>,
411 stream: &ReadableStream,
412 ) -> Fallible<DomRoot<Self>> {
413 let reader = Self::new_with_proto(cx, global, proto);
414
415 reader.set_up(cx, stream, global)?;
417
418 Ok(reader)
419 }
420
421 fn Read(
423 &self,
424 cx: &mut JSContext,
425 view: CustomAutoRooterGuard<ArrayBufferView>,
426 options: &ReadableStreamBYOBReaderReadOptions,
427 ) -> Rc<Promise> {
428 let view = HeapBufferSource::<ArrayBufferViewU8>::from_view(cx, view);
429 let min = options.min;
430 let promise = Promise::new(cx, &self.global());
432
433 if view.byte_length() == 0 {
435 promise.reject_error(cx, Error::Type(c"view byte length is 0".to_owned()));
436 return promise;
437 }
438 if view.viewed_buffer_array_byte_length(cx) == 0 {
441 promise.reject_error(
442 cx,
443 Error::Type(c"viewed buffer byte length is 0".to_owned()),
444 );
445 return promise;
446 }
447
448 if view.is_detached_buffer(cx) {
451 promise.reject_error(cx, Error::Type(c"view is detached".to_owned()));
452 return promise;
453 }
454
455 if min == 0 {
457 promise.reject_error(cx, Error::Type(c"min is 0".to_owned()));
458 return promise;
459 }
460
461 if view.has_typed_array_name() {
463 if min > (view.get_typed_array_length() as u64) {
465 promise.reject_error(
466 cx,
467 Error::Range(c"min is greater than array length".to_owned()),
468 );
469 return promise;
470 }
471 } else {
472 if min > (view.byte_length() as u64) {
475 promise.reject_error(
476 cx,
477 Error::Range(c"min is greater than byte length".to_owned()),
478 );
479 return promise;
480 }
481 }
482
483 if self.stream.get().is_none() {
485 promise.reject_error(
486 cx,
487 Error::Type(c"min is greater than byte length".to_owned()),
488 );
489 return promise;
490 }
491
492 let read_into_request = ReadIntoRequest::Read(promise.clone());
503
504 self.read(cx, &view, min, &read_into_request);
506
507 promise
509 }
510
511 fn ReleaseLock(&self, cx: &mut JSContext) -> Fallible<()> {
513 if self.stream.get().is_none() {
514 return Ok(());
516 }
517
518 self.release(cx)
520 }
521
522 fn Closed(&self) -> Rc<Promise> {
524 self.closed()
525 }
526
527 fn Cancel(&self, cx: &mut JSContext, reason: SafeHandleValue) -> Rc<Promise> {
529 self.generic_cancel(cx, &self.global(), reason)
530 }
531}
532
533impl ReadableStreamGenericReader for ReadableStreamBYOBReader {
534 fn get_closed_promise(&self) -> Rc<Promise> {
535 self.closed_promise.borrow().clone()
536 }
537
538 fn set_closed_promise(&self, promise: Rc<Promise>) {
539 *self.closed_promise.borrow_mut() = promise;
540 }
541
542 fn set_stream(&self, stream: Option<&ReadableStream>) {
543 self.stream.set(stream);
544 }
545
546 fn get_stream(&self) -> Option<DomRoot<ReadableStream>> {
547 self.stream.get()
548 }
549
550 fn as_byob_reader(&self) -> Option<&ReadableStreamBYOBReader> {
551 Some(self)
552 }
553}