script/dom/stream/
byteteereadintorequest.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::jsval::UndefinedValue;
13use js::typedarray::ArrayBufferViewU8;
14use script_bindings::reflector::{Reflector, reflect_dom_object};
15use script_bindings::trace::RootedTraceableBox;
16
17use super::byteteeunderlyingsource::ByteTeePullAlgorithm;
18use crate::dom::bindings::buffer_source::HeapBufferSource;
19use crate::dom::bindings::error::{ErrorToJsval, Fallible};
20use crate::dom::bindings::refcounted::Trusted;
21use crate::dom::bindings::reflector::DomGlobal;
22use crate::dom::bindings::root::{Dom, DomRoot};
23use crate::dom::globalscope::GlobalScope;
24use crate::dom::promise::{RootedPromise, TracedPromise};
25use crate::dom::stream::byteteeunderlyingsource::ByteTeeUnderlyingSource;
26use crate::dom::stream::readablestream::ReadableStream;
27use crate::runtime::job_queue::MicrotaskRunnable;
28
29#[derive(JSTraceable, MallocSizeOf)]
30pub(crate) struct ByteTeeReadIntoRequestMicrotask {
31 #[ignore_malloc_size_of = "mozjs"]
32 chunk: RootedTraceableBox<HeapBufferSource<ArrayBufferViewU8>>,
33 tee_read_request: Trusted<ByteTeeReadIntoRequest>,
34}
35
36impl MicrotaskRunnable for ByteTeeReadIntoRequestMicrotask {
37 fn handler(&self, cx: &mut JSContext) {
38 self.tee_read_request
39 .root()
40 .chunk_steps(&self.chunk, cx)
41 .expect("Failed to enqueue chunk");
42 }
43}
44
45#[dom_struct]
46pub(crate) struct ByteTeeReadIntoRequest {
47 reflector_: Reflector,
48 for_branch2: bool,
49 byob_branch: Dom<ReadableStream>,
50 other_branch: Dom<ReadableStream>,
51 stream: Dom<ReadableStream>,
52 #[conditional_malloc_size_of]
53 read_again_for_branch_1: Rc<Cell<bool>>,
54 #[conditional_malloc_size_of]
55 read_again_for_branch_2: Rc<Cell<bool>>,
56 #[conditional_malloc_size_of]
57 reading: Rc<Cell<bool>>,
58 #[conditional_malloc_size_of]
59 canceled_1: Rc<Cell<bool>>,
60 #[conditional_malloc_size_of]
61 canceled_2: Rc<Cell<bool>>,
62 cancel_promise: TracedPromise,
63 tee_underlying_source: Dom<ByteTeeUnderlyingSource>,
64}
65impl ByteTeeReadIntoRequest {
66 #[allow(clippy::too_many_arguments)]
67 pub(crate) fn new(
68 cx: &mut JSContext,
69 for_branch2: bool,
70 byob_branch: &ReadableStream,
71 other_branch: &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(ByteTeeReadIntoRequest {
85 reflector_: Reflector::new(),
86 for_branch2,
87 byob_branch: Dom::from_ref(byob_branch),
88 other_branch: Dom::from_ref(other_branch),
89 stream: Dom::from_ref(stream),
90 read_again_for_branch_1,
91 read_again_for_branch_2,
92 reading,
93 canceled_1,
94 canceled_2,
95 cancel_promise: cancel_promise.to_traced(),
96 tee_underlying_source: Dom::from_ref(tee_underlying_source),
97 }),
98 global,
99 )
100 }
101
102 pub(crate) fn enqueue_chunk_steps(
103 &self,
104 cx: &mut JSContext,
105 chunk: RootedTraceableBox<HeapBufferSource<ArrayBufferViewU8>>,
106 ) {
107 let byte_tee_read_request_chunk = ByteTeeReadIntoRequestMicrotask {
109 chunk,
110 tee_read_request: Trusted::new(self),
111 };
112
113 self.global()
114 .enqueue_microtask(cx, Box::new(byte_tee_read_request_chunk));
115 }
116
117 #[allow(clippy::borrowed_box)]
119 pub(crate) fn chunk_steps(
120 &self,
121 chunk: &HeapBufferSource<ArrayBufferViewU8>,
122 cx: &mut JSContext,
123 ) -> Fallible<()> {
124 self.read_again_for_branch_1.set(false);
126
127 self.read_again_for_branch_2.set(false);
129
130 let byob_canceled = if self.for_branch2 {
132 self.canceled_2.get()
133 } else {
134 self.canceled_1.get()
135 };
136
137 let other_canceled = if self.for_branch2 {
139 self.canceled_1.get()
140 } else {
141 self.canceled_2.get()
142 };
143
144 if !other_canceled {
146 let clone_result = chunk.clone_as_uint8_array(cx);
148
149 if let Err(error) = clone_result {
151 rooted!(&in(cx) let mut error_value = UndefinedValue());
152 error.to_jsval(cx, &self.global(), error_value.handle_mut());
153
154 let byob_branch_controller = self.byob_branch.get_byte_controller();
156 byob_branch_controller.error(cx, error_value.handle());
157
158 let other_branch_controller = self.other_branch.get_byte_controller();
160 other_branch_controller.error(cx, error_value.handle());
161
162 let cancel_result =
164 self.stream
165 .cancel(cx, &self.stream.global(), error_value.handle());
166 self.cancel_promise.resolve_native(cx, &cancel_result);
167
168 return Ok(());
170 } else {
171 let cloned_chunk = clone_result.unwrap();
173
174 if !byob_canceled {
177 let byob_branch_controller = self.byob_branch.get_byte_controller();
178 byob_branch_controller.respond_with_new_view(cx, chunk)?;
179 }
180
181 let other_branch_controller = self.other_branch.get_byte_controller();
183 other_branch_controller.enqueue(cx, cloned_chunk)?;
184 }
185 } else if !byob_canceled {
186 let byob_branch_controller = self.byob_branch.get_byte_controller();
190 byob_branch_controller.respond_with_new_view(cx, chunk)?;
191 }
192
193 self.reading.set(false);
195
196 if self.read_again_for_branch_1.get() {
198 self.pull_algorithm(cx, Some(ByteTeePullAlgorithm::Pull1Algorithm));
199 } else if self.read_again_for_branch_2.get() {
200 self.pull_algorithm(cx, Some(ByteTeePullAlgorithm::Pull2Algorithm));
202 }
203
204 Ok(())
205 }
206
207 pub(crate) fn close_steps(
209 &self,
210 cx: &mut JSContext,
211 chunk: Option<RootedTraceableBox<HeapBufferSource<ArrayBufferViewU8>>>,
212 ) -> Fallible<()> {
213 self.reading.set(false);
215
216 let byob_canceled = if self.for_branch2 {
218 self.canceled_2.get()
219 } else {
220 self.canceled_1.get()
221 };
222
223 let other_canceled = if self.for_branch2 {
225 self.canceled_1.get()
226 } else {
227 self.canceled_2.get()
228 };
229
230 if !byob_canceled {
232 let byob_branch_controller = self.byob_branch.get_byte_controller();
233 byob_branch_controller.close(cx)?;
234 }
235
236 if !other_canceled {
238 let other_branch_controller = self.other_branch.get_byte_controller();
239 other_branch_controller.close(cx)?;
240 }
241
242 if let Some(chunk_value) = chunk {
244 if chunk_value.is_undefined() {
245 } else {
248 let chunk = chunk_value;
249 assert_eq!(chunk.byte_length(), 0);
251
252 if !byob_canceled {
255 let byob_branch_controller = self.byob_branch.get_byte_controller();
256 byob_branch_controller.respond_with_new_view(cx, &chunk)?;
257 }
258
259 if !other_canceled {
262 let other_branch_controller = self.other_branch.get_byte_controller();
263 if other_branch_controller.get_pending_pull_intos_size() > 0 {
264 other_branch_controller.respond(cx, 0)?;
265 }
266 }
267 }
268 }
269
270 if !byob_canceled || !other_canceled {
272 self.cancel_promise.resolve_native(cx, &());
273 }
274
275 Ok(())
276 }
277 pub(crate) fn error_steps(&self) {
279 self.reading.set(false);
281 }
282
283 pub(crate) fn pull_algorithm(
284 &self,
285 cx: &mut JSContext,
286 byte_tee_pull_algorithm: Option<ByteTeePullAlgorithm>,
287 ) {
288 self.tee_underlying_source
289 .pull_algorithm(cx, byte_tee_pull_algorithm);
290 }
291}