script/dom/stream/
byteteereadrequest.rs1#![cfg_attr(crown, allow(crown::jscontext_first_arg))]
6
7use std::cell::Cell;
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::typedarray::ArrayBufferViewU8;
15use script_bindings::error::Fallible;
16use script_bindings::reflector::{Reflector, reflect_dom_object};
17
18use super::byteteeunderlyingsource::ByteTeePullAlgorithm;
19use crate::dom::bindings::buffer_source::HeapBufferSource;
20use crate::dom::bindings::error::{Error, ErrorToJsval};
21use crate::dom::bindings::reflector::DomGlobal;
22use crate::dom::bindings::root::{Dom, DomRoot};
23use crate::dom::bindings::trace::RootedTraceableBox;
24use crate::dom::globalscope::GlobalScope;
25use crate::dom::promise::{RootedPromise, TracedPromise};
26use crate::dom::stream::byteteeunderlyingsource::ByteTeeUnderlyingSource;
27use crate::dom::stream::readablestream::ReadableStream;
28use crate::runtime::job_queue::MicrotaskRunnable;
29
30#[derive(JSTraceable, MallocSizeOf)]
31#[cfg_attr(crown, expect(crown::unrooted_must_root))]
32pub(crate) struct ByteTeeReadRequestMicrotask {
33 #[ignore_malloc_size_of = "mozjs"]
34 chunk: Box<Heap<JSVal>>,
35 tee_read_request: Dom<ByteTeeReadRequest>,
36}
37
38impl MicrotaskRunnable for ByteTeeReadRequestMicrotask {
39 fn handler(&self, cx: &mut JSContext) {
40 self.tee_read_request
41 .chunk_steps(&self.chunk, cx)
42 .expect("ByteTeeReadRequestMicrotask::microtask_chunk_steps failed");
43 }
44}
45
46#[dom_struct]
47pub(crate) struct ByteTeeReadRequest {
49 reflector_: Reflector,
50 branch_1: Dom<ReadableStream>,
51 branch_2: Dom<ReadableStream>,
52 stream: Dom<ReadableStream>,
53 #[conditional_malloc_size_of]
54 read_again_for_branch_1: Rc<Cell<bool>>,
55 #[conditional_malloc_size_of]
56 read_again_for_branch_2: Rc<Cell<bool>>,
57 #[conditional_malloc_size_of]
58 reading: Rc<Cell<bool>>,
59 #[conditional_malloc_size_of]
60 canceled_1: Rc<Cell<bool>>,
61 #[conditional_malloc_size_of]
62 canceled_2: Rc<Cell<bool>>,
63 cancel_promise: TracedPromise,
64 tee_underlying_source: Dom<ByteTeeUnderlyingSource>,
65}
66impl ByteTeeReadRequest {
67 #[allow(clippy::too_many_arguments)]
68 pub(crate) fn new(
69 cx: &mut JSContext,
70 branch_1: &ReadableStream,
71 branch_2: &ReadableStream,
72 stream: &ReadableStream,
73 read_again_for_branch_1: Rc<Cell<bool>>,
74 read_again_for_branch_2: Rc<Cell<bool>>,
75 reading: Rc<Cell<bool>>,
76 canceled_1: Rc<Cell<bool>>,
77 canceled_2: Rc<Cell<bool>>,
78 cancel_promise: &RootedPromise,
79 tee_underlying_source: &ByteTeeUnderlyingSource,
80 global: &GlobalScope,
81 ) -> DomRoot<Self> {
82 reflect_dom_object(
83 cx,
84 Box::new(ByteTeeReadRequest {
85 reflector_: Reflector::new(),
86 branch_1: Dom::from_ref(branch_1),
87 branch_2: Dom::from_ref(branch_2),
88 stream: Dom::from_ref(stream),
89 read_again_for_branch_1,
90 read_again_for_branch_2,
91 reading,
92 canceled_1,
93 canceled_2,
94 cancel_promise: cancel_promise.to_traced(),
95 tee_underlying_source: Dom::from_ref(tee_underlying_source),
96 }),
97 global,
98 )
99 }
100
101 pub(crate) fn enqueue_chunk_steps(
104 &self,
105 cx: &mut JSContext,
106 global: &GlobalScope,
107 chunk: RootedTraceableBox<Heap<JSVal>>,
108 ) {
109 let byte_tee_read_request_chunk = ByteTeeReadRequestMicrotask {
111 chunk: Heap::boxed(*chunk.handle()),
112 tee_read_request: Dom::from_ref(self),
113 };
114 global.enqueue_microtask(cx, Box::new(byte_tee_read_request_chunk));
115 }
116
117 #[allow(clippy::borrowed_box)]
119 pub(crate) fn chunk_steps(&self, chunk: &Box<Heap<JSVal>>, cx: &mut JSContext) -> Fallible<()> {
120 self.read_again_for_branch_1.set(false);
122
123 self.read_again_for_branch_2.set(false);
125
126 rooted!(&in(cx) let chunk_object = chunk.get().to_object());
128
129 let handle_clone_error = |cx: &mut JSContext, error: Error| {
131 rooted!(&in(cx) let mut error_value = UndefinedValue());
132 error.to_jsval(cx, &self.global(), error_value.handle_mut());
133
134 let branch_1_controller = self.branch_1.get_byte_controller();
135 let branch_2_controller = self.branch_2.get_byte_controller();
136
137 branch_1_controller.error(cx, error_value.handle());
138 branch_2_controller.error(cx, error_value.handle());
139
140 let cancel_result = self
141 .stream
142 .cancel(cx, &self.stream.global(), error_value.handle());
143 self.cancel_promise.resolve_native(cx, &cancel_result);
144 };
145
146 let chunk1_view = if !self.canceled_1.get() {
148 Some(RootedTraceableBox::new(
149 HeapBufferSource::<ArrayBufferViewU8>::new(chunk_object.handle()),
150 ))
151 } else {
152 None
153 };
154
155 let mut chunk2_view = None;
156
157 if !self.canceled_1.get() && !self.canceled_2.get() {
159 let chunk2_source = RootedTraceableBox::new(
161 HeapBufferSource::<ArrayBufferViewU8>::new(chunk_object.handle()),
162 );
163 let clone_result = chunk2_source.clone_as_uint8_array(cx);
164
165 if let Err(error) = clone_result {
167 handle_clone_error(cx, error);
168 return Ok(());
169 } else {
170 chunk2_view = clone_result.ok();
172 }
173 } else if !self.canceled_2.get() {
174 let chunk2_source = RootedTraceableBox::new(
176 HeapBufferSource::<ArrayBufferViewU8>::new(chunk_object.handle()),
177 );
178 match chunk2_source.clone_as_uint8_array(cx) {
179 Ok(clone) => chunk2_view = Some(clone),
180 Err(error) => {
181 handle_clone_error(cx, error);
182 return Ok(());
183 },
184 }
185 }
186
187 if let Some(chunk1_view) = chunk1_view {
189 let branch_1_controller = self.branch_1.get_byte_controller();
190 branch_1_controller.enqueue(cx, chunk1_view)?;
191 }
192
193 if let Some(chunk2_view) = chunk2_view {
195 let branch_2_controller = self.branch_2.get_byte_controller();
196 branch_2_controller.enqueue(cx, chunk2_view)?;
197 }
198
199 self.reading.set(false);
201
202 if self.read_again_for_branch_1.get() {
204 self.pull_algorithm(cx, Some(ByteTeePullAlgorithm::Pull1Algorithm));
205 } else if self.read_again_for_branch_2.get() {
206 self.pull_algorithm(cx, Some(ByteTeePullAlgorithm::Pull2Algorithm));
208 }
209
210 Ok(())
211 }
212
213 pub(crate) fn close_steps(&self, cx: &mut JSContext) -> Fallible<()> {
215 let branch_1_controller = self.branch_1.get_byte_controller();
216 let branch_2_controller = self.branch_2.get_byte_controller();
217
218 self.reading.set(false);
220
221 if !self.canceled_1.get() {
223 branch_1_controller.close(cx)?;
224 }
225
226 if !self.canceled_2.get() {
228 branch_2_controller.close(cx)?;
229 }
230
231 if branch_1_controller.get_pending_pull_intos_size() > 0 {
234 branch_1_controller.respond(cx, 0)?;
235 }
236
237 if branch_2_controller.get_pending_pull_intos_size() > 0 {
240 branch_2_controller.respond(cx, 0)?;
241 }
242
243 if !self.canceled_1.get() || !self.canceled_2.get() {
245 self.cancel_promise.resolve_native(cx, &());
246 }
247
248 Ok(())
249 }
250
251 pub(crate) fn error_steps(&self) {
253 self.reading.set(false);
255 }
256
257 pub(crate) fn pull_algorithm(
258 &self,
259 cx: &mut JSContext,
260 byte_tee_pull_algorithm: Option<ByteTeePullAlgorithm>,
261 ) {
262 self.tee_underlying_source
263 .pull_algorithm(cx, byte_tee_pull_algorithm);
264 }
265}