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, RootedPromise, TracedPromise};
38use crate::dom::promisenativehandler::{Callback, PromiseNativeHandler};
39use crate::dom::stream::readablestream::ReadableStream;
40use crate::realms::enter_auto_realm;
41
42#[derive(Clone, JSTraceable, MallocSizeOf)]
44#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
45pub enum ReadIntoRequest {
46 Read(TracedPromise),
48 ByteTee {
49 byte_tee_read_into_request: Dom<ByteTeeReadIntoRequest>,
50 },
51}
52
53impl js::gc::Rootable for ReadIntoRequest {}
54
55impl ReadIntoRequest {
56 pub fn chunk_steps(&self, cx: &mut JSContext, chunk: RootedTraceableBox<Heap<JSVal>>) {
58 match self {
59 ReadIntoRequest::Read(promise) => {
60 promise.resolve_native(
63 cx,
64 &ReadableStreamReadResult {
65 done: Some(false),
66 value: chunk,
67 },
68 );
69 },
70 ReadIntoRequest::ByteTee {
71 byte_tee_read_into_request,
72 } => {
73 rooted!(&in(cx) let chunk_object = chunk.get().to_object());
74 byte_tee_read_into_request.enqueue_chunk_steps(
75 cx,
76 RootedTraceableBox::new(HeapBufferSource::<ArrayBufferViewU8>::new(
77 chunk_object.handle(),
78 )),
79 )
80 },
81 }
82 }
83
84 pub fn close_steps(&self, cx: &mut JSContext, chunk: Option<RootedTraceableBox<Heap<JSVal>>>) {
86 match self {
87 ReadIntoRequest::Read(promise) => match chunk {
88 Some(chunk) => promise.resolve_native(
91 cx,
92 &ReadableStreamReadResult {
93 done: Some(true),
94 value: chunk,
95 },
96 ),
97 None => {
98 let result = RootedTraceableBox::new(Heap::default());
99 result.set(UndefinedValue());
100 promise.resolve_native(
101 cx,
102 &ReadableStreamReadResult {
103 done: Some(true),
104 value: result,
105 },
106 );
107 },
108 },
109 ReadIntoRequest::ByteTee {
110 byte_tee_read_into_request,
111 } => match chunk {
112 Some(chunk) => {
113 rooted!(&in(cx) let chunk_object = chunk.get().to_object());
114 byte_tee_read_into_request
115 .close_steps(
116 cx,
117 Some(RootedTraceableBox::new(
118 HeapBufferSource::<ArrayBufferViewU8>::new(chunk_object.handle()),
119 )),
120 )
121 .expect("close steps should not fail")
122 },
123 None => byte_tee_read_into_request
124 .close_steps(cx, None)
125 .expect("close steps should not fail"),
126 },
127 }
128 }
129
130 pub(crate) fn error_steps(&self, cx: &mut JSContext, e: SafeHandleValue) {
132 match self {
133 ReadIntoRequest::Read(promise) => {
134 promise.reject_native(cx, &e)
137 },
138 ReadIntoRequest::ByteTee {
139 byte_tee_read_into_request,
140 } => {
141 byte_tee_read_into_request.error_steps();
142 },
143 }
144 }
145}
146
147#[derive(Clone, JSTraceable, MallocSizeOf)]
150#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
151struct ByteTeeClosedPromiseRejectionHandler {
152 branch_1_controller: Dom<ReadableByteStreamController>,
153 branch_2_controller: Dom<ReadableByteStreamController>,
154 #[conditional_malloc_size_of]
155 canceled_1: Rc<Cell<bool>>,
156 #[conditional_malloc_size_of]
157 canceled_2: Rc<Cell<bool>>,
158 cancel_promise: TracedPromise,
159 #[conditional_malloc_size_of]
160 reader_version: Rc<Cell<u64>>,
161 expected_version: u64,
162}
163
164impl Callback for ByteTeeClosedPromiseRejectionHandler {
165 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
168 if self.reader_version.get() != self.expected_version {
170 return;
171 }
172
173 self.branch_1_controller.error(cx, v);
175
176 self.branch_2_controller.error(cx, v);
178
179 if !self.canceled_1.get() || !self.canceled_2.get() {
181 self.cancel_promise.resolve_native(cx, &());
182 }
183 }
184}
185
186#[dom_struct]
188pub(crate) struct ReadableStreamBYOBReader {
189 reflector_: Reflector,
190
191 stream: MutNullableDom<ReadableStream>,
193
194 read_into_requests: DomRefCell<VecDeque<ReadIntoRequest>>,
195
196 closed_promise: DomRefCell<TracedPromise>,
198}
199
200impl ReadableStreamBYOBReader {
201 fn new_with_proto(
202 cx: &mut JSContext,
203 global: &GlobalScope,
204 proto: Option<SafeHandleObject>,
205 ) -> DomRoot<ReadableStreamBYOBReader> {
206 let closed_promise = Promise::new_rooted(cx, global);
207 reflect_dom_object_with_proto(
208 cx,
209 Box::new(ReadableStreamBYOBReader::new_inherited(&closed_promise)),
210 global,
211 proto,
212 )
213 }
214
215 fn new_inherited(promise: &RootedPromise) -> ReadableStreamBYOBReader {
216 ReadableStreamBYOBReader {
217 reflector_: Reflector::new(),
218 stream: MutNullableDom::new(None),
219 read_into_requests: DomRefCell::new(Default::default()),
220 closed_promise: DomRefCell::new(promise.to_traced()),
221 }
222 }
223
224 pub(crate) fn new(
225 cx: &mut JSContext,
226 global: &GlobalScope,
227 ) -> DomRoot<ReadableStreamBYOBReader> {
228 let closed_promise = Promise::new_rooted(cx, global);
229 reflect_dom_object_with_cx(Box::new(Self::new_inherited(&closed_promise)), global, cx)
230 }
231
232 pub(crate) fn set_up(
234 &self,
235 cx: &mut JSContext,
236 stream: &ReadableStream,
237 global: &GlobalScope,
238 ) -> Fallible<()> {
239 if stream.is_locked() {
241 return Err(Error::Type(c"stream is locked".to_owned()));
242 }
243
244 if !stream.has_byte_controller() {
246 return Err(Error::Type(
247 c"stream controller is not a byte stream controller".to_owned(),
248 ));
249 }
250
251 self.generic_initialize(cx, global, stream);
253
254 self.read_into_requests.borrow_mut().clear();
256
257 Ok(())
258 }
259
260 pub(crate) fn release(&self, cx: &mut JSContext) -> Fallible<()> {
262 self.generic_release(cx).expect("Generic release failed");
264 rooted!(&in(cx) let mut error = UndefinedValue());
266 Error::Type(c"Reader is released".to_owned()).to_jsval(
267 cx,
268 &self.global(),
269 error.handle_mut(),
270 );
271
272 self.error_read_into_requests(cx, error.handle());
274 Ok(())
275 }
276
277 pub(crate) fn error_read_into_requests(&self, cx: &mut JSContext, e: SafeHandleValue) {
279 self.closed_promise.borrow().reject_native(cx, &e);
281
282 self.closed_promise.borrow().set_promise_is_handled(cx);
284
285 rooted!(&in(cx) let read_into_requests = self.take_read_into_requests());
287
288 for request in read_into_requests.iter() {
290 request.error_steps(cx, e);
292 }
293 }
294
295 fn take_read_into_requests(&self) -> VecDeque<ReadIntoRequest> {
296 mem::take(&mut *self.read_into_requests.borrow_mut())
297 }
298
299 pub(crate) fn add_read_into_request(&self, read_request: &ReadIntoRequest) {
301 self.read_into_requests
302 .borrow_mut()
303 .push_back(read_request.clone());
304 }
305
306 pub(crate) fn cancel(&self, cx: &mut JSContext) {
308 rooted!(&in(cx) let read_into_requests = self.take_read_into_requests());
311 for request in read_into_requests.iter() {
314 request.close_steps(cx, None);
316 }
317 }
318
319 pub(crate) fn close(&self, cx: &mut JSContext) {
320 self.closed_promise.borrow().resolve_native(cx, &());
322 }
323
324 pub(crate) fn read(
326 &self,
327 cx: &mut JSContext,
328 view: &HeapBufferSource<ArrayBufferViewU8>,
329 min: u64,
330 read_into_request: &ReadIntoRequest,
331 ) {
332 assert!(self.stream.get().is_some());
336
337 let stream = self.stream.get().unwrap();
338
339 stream.set_is_disturbed(true);
341 if stream.is_errored() {
343 rooted!(&in(cx) let mut error = UndefinedValue());
344 stream.get_stored_error(error.handle_mut());
345
346 read_into_request.error_steps(cx, error.handle());
347 } else {
348 stream.perform_pull_into(cx, read_into_request, view, min);
351 }
352 }
353
354 pub(crate) fn get_num_read_into_requests(&self) -> usize {
355 self.read_into_requests.borrow().len()
356 }
357
358 pub(crate) fn remove_read_into_request(&self) -> ReadIntoRequest {
359 self.read_into_requests
360 .borrow_mut()
361 .pop_front()
362 .expect("read into requests is empty")
363 }
364
365 #[allow(clippy::too_many_arguments)]
366 pub(crate) fn byte_tee_append_native_handler_to_closed_promise(
367 &self,
368 cx: &mut JSContext,
369 branch_1: &ReadableStream,
370 branch_2: &ReadableStream,
371 canceled_1: Rc<Cell<bool>>,
372 canceled_2: Rc<Cell<bool>>,
373 cancel_promise: &RootedPromise,
374 reader_version: Rc<Cell<u64>>,
375 expected_version: u64,
376 ) {
377 let branch_1_controller = branch_1.get_byte_controller();
378
379 let branch_2_controller = branch_2.get_byte_controller();
380
381 let global = self.global();
382 let handler = PromiseNativeHandler::new(
383 cx,
384 &global,
385 None,
386 Some(Box::new(ByteTeeClosedPromiseRejectionHandler {
387 branch_1_controller: Dom::from_ref(&branch_1_controller),
388 branch_2_controller: Dom::from_ref(&branch_2_controller),
389 canceled_1,
390 canceled_2,
391 cancel_promise: cancel_promise.to_traced(),
392 reader_version,
393 expected_version,
394 })),
395 );
396
397 let mut realm = enter_auto_realm(cx, &*global);
398 let cx = &mut realm.current_realm();
399
400 self.closed_promise
401 .borrow()
402 .append_native_handler(cx, &handler);
403 }
404}
405
406impl ReadableStreamBYOBReaderMethods<crate::DomTypeHolder> for ReadableStreamBYOBReader {
407 fn Constructor(
409 cx: &mut JSContext,
410 global: &GlobalScope,
411 proto: Option<SafeHandleObject>,
412 stream: &ReadableStream,
413 ) -> Fallible<DomRoot<Self>> {
414 let reader = Self::new_with_proto(cx, global, proto);
415
416 reader.set_up(cx, stream, global)?;
418
419 Ok(reader)
420 }
421
422 fn Read(
424 &self,
425 cx: &mut JSContext,
426 view: CustomAutoRooterGuard<ArrayBufferView>,
427 options: &ReadableStreamBYOBReaderReadOptions,
428 ) -> RootedPromise {
429 let view = HeapBufferSource::<ArrayBufferViewU8>::from_view(cx, view);
430 let min = options.min;
431 let promise = Promise::new_rooted(cx, &self.global());
433
434 if view.byte_length() == 0 {
436 promise.reject_error(cx, Error::Type(c"view byte length is 0".to_owned()));
437 return promise;
438 }
439 if view.viewed_buffer_array_byte_length(cx) == 0 {
442 promise.reject_error(
443 cx,
444 Error::Type(c"viewed buffer byte length is 0".to_owned()),
445 );
446 return promise;
447 }
448
449 if view.is_detached_buffer(cx) {
452 promise.reject_error(cx, Error::Type(c"view is detached".to_owned()));
453 return promise;
454 }
455
456 if min == 0 {
458 promise.reject_error(cx, Error::Type(c"min is 0".to_owned()));
459 return promise;
460 }
461
462 if view.has_typed_array_name() {
464 if min > (view.get_typed_array_length() as u64) {
466 promise.reject_error(
467 cx,
468 Error::Range(c"min is greater than array length".to_owned()),
469 );
470 return promise;
471 }
472 } else {
473 if min > (view.byte_length() as u64) {
476 promise.reject_error(
477 cx,
478 Error::Range(c"min is greater than byte length".to_owned()),
479 );
480 return promise;
481 }
482 }
483
484 if self.stream.get().is_none() {
486 promise.reject_error(
487 cx,
488 Error::Type(c"min is greater than byte length".to_owned()),
489 );
490 return promise;
491 }
492
493 rooted!(&in(cx) let read_into_request = ReadIntoRequest::Read(promise.to_traced()));
504
505 self.read(cx, &view, min, &read_into_request);
507
508 promise
510 }
511
512 fn ReleaseLock(&self, cx: &mut JSContext) -> Fallible<()> {
514 if self.stream.get().is_none() {
515 return Ok(());
517 }
518
519 self.release(cx)
521 }
522
523 fn Closed(&self, cx: &JSContext) -> RootedPromise {
525 self.closed(cx)
526 }
527
528 fn Cancel(&self, cx: &mut JSContext, reason: SafeHandleValue) -> RootedPromise {
530 self.generic_cancel(cx, &self.global(), reason)
531 }
532}
533
534impl ReadableStreamGenericReader for ReadableStreamBYOBReader {
535 fn get_closed_promise(&self, cx: &JSContext) -> RootedPromise {
536 self.closed_promise.borrow().root(cx)
537 }
538
539 fn set_closed_promise(&self, promise: &RootedPromise) {
540 *self.closed_promise.borrow_mut() = promise.to_traced();
541 }
542
543 fn set_stream(&self, stream: Option<&ReadableStream>) {
544 self.stream.set(stream);
545 }
546
547 fn get_stream(&self) -> Option<DomRoot<ReadableStream>> {
548 self.stream.get()
549 }
550
551 fn as_byob_reader(&self) -> Option<&ReadableStreamBYOBReader> {
552 Some(self)
553 }
554}