use std::collections::VecDeque;
use crate::utils::{bytes::ObjectBytes, primordials::Primordial};
use rquickjs::{
atom::PredefinedAtom,
class::{JsClass, OwnedBorrowMut, Trace, Tracer},
function::Constructor,
methods,
prelude::{Opt, This},
ArrayBuffer, Class, Ctx, Error, Exception, FromJs, Function, IntoJs, JsLifetime, Object,
Promise, Result, Value,
};
use crate::stream_web::{
readable::{
byte_controller::ReadableByteStreamController,
controller::{ReadableStreamController, ReadableStreamControllerClass},
default_reader::{ReadableStreamDefaultReaderOwned, ReadableStreamReadResult},
objects::{ReadableStreamBYOBObjects, ReadableStreamObjects},
reader::{ReadableStreamGenericReader, ReadableStreamReader, ReadableStreamReaderOwned},
stream::{ReadableStreamOwned, ReadableStreamState},
},
utils::{
promise::{promise_rejected_with_constructor, with_promise_result, ResolveablePromise},
UnwrapOrUndefined, ValueOrUndefined,
},
};
#[derive(Trace)]
#[rquickjs::class]
pub(crate) struct ReadableStreamBYOBReader<'js> {
pub(super) generic: ReadableStreamGenericReader<'js>,
pub(super) read_into_requests: VecDeque<Box<dyn ReadableStreamReadIntoRequest<'js> + 'js>>,
}
pub(crate) type ReadableStreamBYOBReaderClass<'js> = Class<'js, ReadableStreamBYOBReader<'js>>;
pub(crate) type ReadableStreamBYOBReaderOwned<'js> =
OwnedBorrowMut<'js, ReadableStreamBYOBReader<'js>>;
unsafe impl<'js> JsLifetime<'js> for ReadableStreamBYOBReader<'js> {
type Changed<'to> = ReadableStreamBYOBReader<'to>;
}
impl<'js> ReadableStreamBYOBReader<'js> {
pub(super) fn readable_stream_byob_reader_error_read_into_requests(
mut objects: ReadableStreamBYOBObjects<'js>,
e: Value<'js>,
) -> Result<ReadableStreamBYOBObjects<'js>> {
let read_into_requests = &mut objects.reader.read_into_requests;
let read_into_requests = read_into_requests.split_off(0);
for read_into_request in read_into_requests {
objects = read_into_request.error_steps(objects, e.clone())?;
}
Ok(objects)
}
pub(super) fn set_up_readable_stream_byob_reader(
ctx: Ctx<'js>,
stream: ReadableStreamOwned<'js>,
) -> Result<(ReadableStreamOwned<'js>, Class<'js, Self>)> {
if stream.is_readable_stream_locked() {
return Err(Exception::throw_type(
&ctx,
"This stream has already been locked for exclusive reading by another reader",
));
}
match stream.controller {
ReadableStreamControllerClass::ReadableStreamByteController(_) => {},
_ => {
return Err(Exception::throw_type(
&ctx,
"Cannot construct a ReadableStreamBYOBReader for a stream not constructed with a byte source",
));
},
};
let generic =
ReadableStreamGenericReader::readable_stream_reader_generic_initialize(&ctx, stream)?;
let mut stream = OwnedBorrowMut::from_class(generic.stream.clone().unwrap());
let reader = Class::instance(
ctx.clone(),
Self {
generic,
read_into_requests: VecDeque::new(),
},
)?;
stream.reader = Some(reader.clone().into());
Ok((stream, reader))
}
pub(super) fn readable_stream_byob_reader_release(
mut objects: ReadableStreamBYOBObjects<'js>,
) -> Result<ReadableStreamBYOBObjects<'js>> {
objects
.reader
.generic
.readable_stream_reader_generic_release(&mut objects.stream, || {
objects.controller.release_steps()
})?;
let e: Value = objects
.stream
.constructor_type_error
.call(("Reader was released",))?;
Self::readable_stream_byob_reader_error_read_into_requests(objects, e)
}
pub(super) fn readable_stream_byob_reader_read(
ctx: &Ctx<'js>,
mut objects: ReadableStreamBYOBObjects<'js>,
view: ViewBytes<'js>,
min: u64,
read_into_request: impl ReadableStreamReadIntoRequest<'js> + 'js,
) -> Result<ReadableStreamBYOBObjects<'js>> {
objects.stream.disturbed = true;
if let ReadableStreamState::Errored(ref stored_error) = objects.stream.state {
let stored_error = stored_error.clone();
read_into_request.error_steps(objects, stored_error)
} else {
ReadableByteStreamController::readable_byte_stream_controller_pull_into(
ctx,
objects,
view,
min,
read_into_request,
)
}
}
}
#[methods(rename_all = "camelCase")]
impl<'js> ReadableStreamBYOBReader<'js> {
#[qjs(get)]
pub fn constructor(ctx: Ctx<'js>) -> Result<Option<Constructor<'js>>> {
<ReadableStreamBYOBReader as JsClass>::constructor(&ctx)
}
#[qjs(constructor)]
pub fn new(ctx: Ctx<'js>, stream: ReadableStreamOwned<'js>) -> Result<Class<'js, Self>> {
let (_, reader) = Self::set_up_readable_stream_byob_reader(ctx, stream)?;
Ok(reader)
}
fn read(
ctx: Ctx<'js>,
reader: This<OwnedBorrowMut<'js, Self>>,
view: Opt<Value<'js>>,
options: Opt<Value<'js>>,
) -> Result<Promise<'js>> {
with_promise_result(&ctx, || {
let options = match options.0 {
None => ReadableStreamBYOBReaderReadOptions { min: 1 },
Some(value) => ReadableStreamBYOBReaderReadOptions::from_js(&ctx, value)?,
};
let view = ViewBytes::from_value(
&ctx,
&reader.generic.function_array_buffer_is_view,
view.0.as_ref(),
)?;
let (buffer, byte_length, _) = view.get_array_buffer()?;
if byte_length == 0 {
return promise_rejected_with_constructor(
&reader.generic.constructor_type_error,
&reader.generic.promise_primordials,
"view must have non-zero byteLength",
);
}
if buffer.is_empty() {
return promise_rejected_with_constructor(
&reader.generic.constructor_type_error,
&reader.generic.promise_primordials,
"view's buffer must have non-zero byteLength",
);
}
if unsafe { buffer.as_bytes() }.is_none() {
return promise_rejected_with_constructor(
&reader.generic.constructor_type_error,
&reader.generic.promise_primordials,
"view's buffer has been detached",
);
}
if options.min == 0 {
return promise_rejected_with_constructor(
&reader.generic.constructor_type_error,
&reader.generic.promise_primordials,
"options.min must be greater than 0",
);
}
let typed_array_len = match &view.0 {
ObjectBytes::U8Array(a) => Some(a.len()),
ObjectBytes::I8Array(a) => Some(a.len()),
ObjectBytes::U16Array(a) => Some(a.len()),
ObjectBytes::I16Array(a) => Some(a.len()),
ObjectBytes::U32Array(a) => Some(a.len()),
ObjectBytes::I32Array(a) => Some(a.len()),
ObjectBytes::U64Array(a) => Some(a.len()),
ObjectBytes::I64Array(a) => Some(a.len()),
ObjectBytes::F32Array(a) => Some(a.len()),
ObjectBytes::F64Array(a) => Some(a.len()),
_ => None,
};
if let Some(typed_array_len) = typed_array_len {
if options.min > typed_array_len as u64 {
return promise_rejected_with_constructor(
&reader.generic.constructor_range_error,
&reader.generic.promise_primordials,
"options.min must be less than or equal to views length",
);
}
} else {
if options.min > byte_length as u64 {
return promise_rejected_with_constructor(
&reader.generic.constructor_range_error,
&reader.generic.promise_primordials,
"options.min must be less than or equal to views byteLength",
);
}
}
if reader.generic.stream.is_none() {
return promise_rejected_with_constructor(
&reader.generic.constructor_type_error,
&reader.generic.promise_primordials,
"Cannot read a stream using a released reader",
);
}
let promise = ResolveablePromise::new(&ctx)?;
#[derive(Trace)]
struct ReadIntoRequest<'js> {
promise: ResolveablePromise<'js>,
}
impl<'js> ReadableStreamReadIntoRequest<'js> for ReadIntoRequest<'js> {
fn chunk_steps(
&self,
objects: ReadableStreamBYOBObjects<'js>,
chunk: Value<'js>,
) -> Result<ReadableStreamBYOBObjects<'js>> {
self.promise.resolve(ReadableStreamReadResult {
value: Some(chunk),
done: false,
})?;
Ok(objects)
}
fn close_steps(
&self,
objects: ReadableStreamBYOBObjects<'js>,
chunk: Value<'js>,
) -> Result<ReadableStreamBYOBObjects<'js>> {
self.promise.resolve(ReadableStreamReadResult {
value: Some(chunk),
done: true,
})?;
Ok(objects)
}
fn error_steps(
&self,
objects: ReadableStreamBYOBObjects<'js>,
reason: Value<'js>,
) -> Result<ReadableStreamBYOBObjects<'js>> {
self.promise.reject(reason)?;
Ok(objects)
}
}
let objects = ReadableStreamObjects::from_byob_reader(reader.0);
Self::readable_stream_byob_reader_read(
&ctx,
objects,
view,
options.min,
ReadIntoRequest {
promise: promise.clone(),
},
)?;
Ok(promise.promise)
})
}
fn release_lock(reader: This<OwnedBorrowMut<'js, Self>>) -> Result<()> {
if reader.generic.stream.is_none() {
return Ok(());
};
let objects = ReadableStreamObjects::from_byob_reader(reader.0);
Self::readable_stream_byob_reader_release(objects)?;
Ok(())
}
#[qjs(get)]
fn closed(&self) -> Promise<'js> {
self.generic.closed_promise.promise.clone()
}
fn cancel(
ctx: Ctx<'js>,
reader: This<OwnedBorrowMut<'js, Self>>,
reason: Opt<Value<'js>>,
) -> Result<Promise<'js>> {
if reader.generic.stream.is_none() {
return promise_rejected_with_constructor(
&reader.generic.constructor_type_error,
&reader.generic.promise_primordials,
"Cannot cancel a stream using a released reader",
);
}
let objects = ReadableStreamObjects::from_byob_reader(reader.0);
let (promise, _) = ReadableStreamGenericReader::readable_stream_reader_generic_cancel(
ctx.clone(),
objects,
reason.0.unwrap_or_undefined(&ctx),
)?;
Ok(promise)
}
}
struct ReadableStreamBYOBReaderReadOptions {
min: u64,
}
impl<'js> FromJs<'js> for ReadableStreamBYOBReaderReadOptions {
fn from_js(ctx: &Ctx<'js>, value: Value<'js>) -> Result<Self> {
let ty_name = value.type_name();
let obj = value
.as_object()
.ok_or(Error::new_from_js(ty_name, "Object"))?;
let min = obj.get_value_or_undefined::<_, f64>("min")?.unwrap_or(1.0);
if min < u64::MIN as f64 || min > u64::MAX as f64 {
return Err(Exception::throw_type(
ctx,
"min on ReadableStreamBYOBReaderReadOptions must fit into unsigned long long",
));
};
Ok(Self { min: min as u64 })
}
}
pub(super) trait ReadableStreamReadIntoRequest<'js>: Trace<'js> {
fn chunk_steps(
&self,
objects: ReadableStreamBYOBObjects<'js>,
chunk: Value<'js>,
) -> Result<ReadableStreamBYOBObjects<'js>>;
fn close_steps(
&self,
objects: ReadableStreamBYOBObjects<'js>,
chunk: Value<'js>,
) -> Result<ReadableStreamBYOBObjects<'js>>;
fn error_steps(
&self,
objects: ReadableStreamBYOBObjects<'js>,
reason: Value<'js>,
) -> Result<ReadableStreamBYOBObjects<'js>>;
}
impl<'js> Trace<'js> for Box<dyn ReadableStreamReadIntoRequest<'js> + 'js> {
fn trace<'a>(&self, tracer: Tracer<'a, 'js>) {
self.as_ref().trace(tracer);
}
}
#[derive(JsLifetime, Clone)]
pub(super) struct ViewBytes<'js>(ObjectBytes<'js>);
impl<'js> ViewBytes<'js> {
pub(super) fn from_object(
ctx: &Ctx<'js>,
function_array_buffer_is_view: &Function<'js>,
object: &Object<'js>,
) -> Result<Self> {
if function_array_buffer_is_view.call::<_, bool>((object.clone(),))? {
if let Some(view) = ObjectBytes::from_array_buffer(object)? {
return Ok(Self(view));
}
}
Err(Exception::throw_type(
ctx,
"view must be an ArrayBufferView",
))
}
pub(super) fn from_value(
ctx: &Ctx<'js>,
function_array_buffer_is_view: &Function<'js>,
value: Option<&Value<'js>>,
) -> Result<Self> {
match value.and_then(Value::as_object) {
None => {
Err(Exception::throw_type(
ctx,
"view must be typed DataView, Buffer, ArrayBuffer, or Uint8Array, but is not an object",
))
},
Some(object) => Self::from_object(ctx, function_array_buffer_is_view, object),
}
}
pub(super) fn get_array_buffer(&self) -> Result<(ArrayBuffer<'js>, usize, usize)> {
Ok(self
.0
.get_array_buffer()?
.expect("invariant broken; ViewBytes may not contain ObjectBytes::Vec"))
}
pub(super) fn element_size(&self) -> usize {
match self.0 {
ObjectBytes::U8Array(_) => 1,
ObjectBytes::I8Array(_) => 1,
ObjectBytes::U16Array(_) => 2,
ObjectBytes::I16Array(_) => 2,
ObjectBytes::U32Array(_) => 4,
ObjectBytes::I32Array(_) => 4,
ObjectBytes::U64Array(_) => 8,
ObjectBytes::I64Array(_) => 8,
ObjectBytes::F16Array(_) => 2,
ObjectBytes::F32Array(_) => 4,
ObjectBytes::F64Array(_) => 8,
ObjectBytes::U8ClampedArray(_) => 1,
ObjectBytes::DataView(_, _, _) => 1,
ObjectBytes::Vec(_) => {
panic!("invariant broken; ViewBytes may not contain ObjectBytes::Vec")
},
}
}
}
#[derive(Clone, JsLifetime)]
pub(crate) struct ArrayConstructorPrimordials<'js> {
pub(super) constructor_uint8array: Constructor<'js>,
constructor_int8array: Constructor<'js>,
constructor_uint16array: Constructor<'js>,
constructor_int16array: Constructor<'js>,
constructor_uint32array: Constructor<'js>,
constructor_int32array: Constructor<'js>,
constructor_uint64array: Constructor<'js>,
constructor_int64array: Constructor<'js>,
constructor_f16array: Constructor<'js>,
constructor_f32array: Constructor<'js>,
constructor_f64array: Constructor<'js>,
constructor_uint8clampedarray: Constructor<'js>,
constructor_data_view: Constructor<'js>,
}
impl<'js> Trace<'js> for ArrayConstructorPrimordials<'js> {
fn trace<'a>(&self, tracer: Tracer<'a, 'js>) {
self.constructor_uint8array.trace(tracer);
self.constructor_int8array.trace(tracer);
self.constructor_uint16array.trace(tracer);
self.constructor_int16array.trace(tracer);
self.constructor_uint32array.trace(tracer);
self.constructor_int32array.trace(tracer);
self.constructor_uint64array.trace(tracer);
self.constructor_int64array.trace(tracer);
self.constructor_f16array.trace(tracer);
self.constructor_f32array.trace(tracer);
self.constructor_f64array.trace(tracer);
self.constructor_uint8clampedarray.trace(tracer);
self.constructor_data_view.trace(tracer);
}
}
impl<'js> Primordial<'js> for ArrayConstructorPrimordials<'js> {
fn new(ctx: &Ctx<'js>) -> Result<Self>
where
Self: Sized,
{
let globals = ctx.globals();
Ok(Self {
constructor_uint8array: globals.get(PredefinedAtom::Uint8Array)?,
constructor_int8array: globals.get(PredefinedAtom::Int8Array)?,
constructor_uint16array: globals.get(PredefinedAtom::Uint16Array)?,
constructor_int16array: globals.get(PredefinedAtom::Int16Array)?,
constructor_uint32array: globals.get(PredefinedAtom::Uint32Array)?,
constructor_int32array: globals.get(PredefinedAtom::Int32Array)?,
constructor_uint64array: globals.get(PredefinedAtom::BigUint64Array)?,
constructor_int64array: globals.get(PredefinedAtom::BigInt64Array)?,
constructor_f16array: globals.get(PredefinedAtom::Float16Array)?,
constructor_f32array: globals.get(PredefinedAtom::Float32Array)?,
constructor_f64array: globals.get(PredefinedAtom::Float64Array)?,
constructor_uint8clampedarray: globals.get(PredefinedAtom::Uint8ClampedArray)?,
constructor_data_view: globals.get(PredefinedAtom::DataView)?,
})
}
}
impl<'js> ArrayConstructorPrimordials<'js> {
pub(super) fn for_view_bytes(&self, v: &ViewBytes<'js>) -> Constructor<'js> {
match v.0 {
ObjectBytes::U8Array(_) => self.constructor_uint8array.clone(),
ObjectBytes::I8Array(_) => self.constructor_int8array.clone(),
ObjectBytes::U16Array(_) => self.constructor_uint16array.clone(),
ObjectBytes::I16Array(_) => self.constructor_int16array.clone(),
ObjectBytes::U32Array(_) => self.constructor_uint32array.clone(),
ObjectBytes::I32Array(_) => self.constructor_int32array.clone(),
ObjectBytes::U64Array(_) => self.constructor_uint64array.clone(),
ObjectBytes::I64Array(_) => self.constructor_int64array.clone(),
ObjectBytes::F16Array(_) => self.constructor_f16array.clone(),
ObjectBytes::F32Array(_) => self.constructor_f32array.clone(),
ObjectBytes::F64Array(_) => self.constructor_f64array.clone(),
ObjectBytes::U8ClampedArray(_) => self.constructor_uint8clampedarray.clone(),
ObjectBytes::DataView(_, _, _) => self.constructor_data_view.clone(),
ObjectBytes::Vec(_) => {
panic!("invariant broken; ViewBytes may not contain ObjectBytes::Vec")
},
}
}
}
impl<'js> Trace<'js> for ViewBytes<'js> {
fn trace<'a>(&self, tracer: Tracer<'a, 'js>) {
self.0.trace(tracer);
}
}
impl<'js> IntoJs<'js> for ViewBytes<'js> {
fn into_js(self, ctx: &Ctx<'js>) -> Result<Value<'js>> {
self.0.into_js(ctx)
}
}
impl<'js> ReadableStreamReader<'js> for ReadableStreamBYOBReaderOwned<'js> {
type Class = ReadableStreamBYOBReaderClass<'js>;
fn with_reader<C>(
self,
ctx: C,
_: impl FnOnce(
C,
ReadableStreamDefaultReaderOwned<'js>,
) -> Result<(C, ReadableStreamDefaultReaderOwned<'js>)>,
byob: impl FnOnce(
C,
ReadableStreamBYOBReaderOwned<'js>,
) -> Result<(C, ReadableStreamBYOBReaderOwned<'js>)>,
_: impl FnOnce(C) -> Result<C>,
) -> Result<(C, Self)> {
byob(ctx, self)
}
fn into_inner(self) -> Self::Class {
self.into_inner()
}
fn from_class(class: Self::Class) -> Self {
OwnedBorrowMut::from_class(class)
}
fn try_from_erased(erased: Option<ReadableStreamReaderOwned<'js>>) -> Option<Self> {
match erased {
Some(ReadableStreamReaderOwned::ReadableStreamBYOBReader(r)) => Some(r),
_ => None,
}
}
}