use std::{
rc::Rc,
sync::atomic::{AtomicBool, Ordering},
};
use crate::utils::primordials::{BasePrimordials, Primordial};
use rquickjs::{
atom::PredefinedAtom,
class::{
impl_::{CloneTrait, CloneWrapper},
JsClass, OwnedBorrow, OwnedBorrowMut, Trace, Tracer,
},
function::Constructor,
methods,
prelude::{Opt, This},
Class, Coerced, Ctx, Error, Exception, FromJs, Function, IntoAtom, IntoJs, JsLifetime, Object,
Promise, Result, Symbol, Type, Value,
};
use crate::stream_web::{
readable::{
controller::ReadableStreamControllerOwned,
default_reader::{
ReadableStreamDefaultReader, ReadableStreamDefaultReaderOwned,
ReadableStreamReadRequest, ReadableStreamReadResult,
},
objects::{
ReadableStreamClassObjects, ReadableStreamDefaultReaderObjects, ReadableStreamObjects,
},
reader::ReadableStreamGenericReader,
},
utils::{
class_from_owned_borrow_mut,
promise::{promise_resolved_with, PromisePrimordials},
promise::{upon_promise, upon_promise_fulfilment, ResolveablePromise},
UnwrapOrUndefined,
},
};
pub(super) enum IteratorKind {
Async,
}
#[derive(Trace)]
pub(super) struct IteratorRecord<'js> {
pub(super) iterator: Object<'js>,
next_method: Function<'js>,
#[qjs(skip_trace)]
done: AtomicBool,
sync_to_async_iterator: Function<'js>,
}
impl<'js> IteratorRecord<'js> {
pub(super) fn get_iterator(
ctx: &Ctx<'js>,
obj: Value<'js>,
kind: IteratorKind,
) -> Result<Self> {
let method: Option<Function<'js>> = match kind {
IteratorKind::Async => {
let method = get_method(ctx, obj.clone(), Symbol::async_iterator(ctx.clone()))?;
if method.is_none() {
let sync_method = get_method(ctx, obj.clone(), Symbol::iterator(ctx.clone()))?;
let sync_method = match sync_method {
None => {
return Err(Exception::throw_type(ctx, "Object is not an iterator"));
},
Some(sync_method) => sync_method,
};
let sync_iterator_record =
Self::get_iterator_from_method(ctx, &obj, sync_method)?;
return sync_iterator_record.create_async_from_sync_iterator(ctx);
}
method
},
};
match method {
None => Err(Exception::throw_type(ctx, "Object is not an iterator")),
Some(method) => {
Self::get_iterator_from_method(ctx, &obj, method)
},
}
}
fn get_iterator_from_method(
ctx: &Ctx<'js>,
obj: &Value<'js>,
method: Function<'js>,
) -> Result<Self> {
let iterator: Value<'js> = method.call((This(obj),))?;
let iterator = match iterator.into_object() {
Some(iterator) => iterator,
None => {
return Err(Exception::throw_type(
ctx,
"The iterator method must return an object",
));
},
};
let next_method = iterator.get(PredefinedAtom::Next)?;
Ok(Self {
iterator,
next_method,
done: AtomicBool::new(false),
sync_to_async_iterator: IteratorPrimordials::get(ctx)?
.sync_to_async_iterator
.clone(),
})
}
fn create_async_from_sync_iterator(self, ctx: &Ctx<'js>) -> Result<Self> {
let sync_iterable = Object::new(ctx.clone())?;
sync_iterable.set(
Symbol::iterator(ctx.clone()),
Function::new(ctx.clone(), {
let iterator = self.iterator.clone();
move || iterator.clone()
}),
)?;
let async_iterator: Object<'js> = self.sync_to_async_iterator.call((sync_iterable,))?;
let next_method = async_iterator.get(PredefinedAtom::Next)?;
Ok(Self {
iterator: async_iterator,
next_method,
done: AtomicBool::new(false),
sync_to_async_iterator: self.sync_to_async_iterator,
})
}
pub(super) fn iterator_next(
&self,
ctx: &Ctx<'js>,
value: Option<Value<'js>>,
) -> Result<Object<'js>> {
let result: Result<Value<'js>> = match value {
None => {
self.next_method.call((This(self.iterator.clone()),))
},
Some(value) => {
self.next_method.call((This(self.iterator.clone()), value))
},
};
let result = match result {
Err(Error::Exception) => {
self.done.store(true, Ordering::Release);
return Err(Error::Exception);
},
Err(err) => return Err(err),
Ok(result) => result,
};
let result = match result.into_object() {
None => {
self.done.store(true, Ordering::Release);
return Err(Exception::throw_type(
ctx,
"The iterator.next() method must return an object",
));
},
Some(result) => result,
};
Ok(result)
}
pub(super) fn iterator_complete(iterator_result: &Object<'js>) -> Result<bool> {
let done: Coerced<bool> = iterator_result.get(PredefinedAtom::Done)?;
Ok(done.0)
}
pub(super) fn iterator_value(iterator_result: &Object<'js>) -> Result<Value<'js>> {
iterator_result.get(PredefinedAtom::Value)
}
}
pub(super) struct ReadableStreamAsyncIterator<'js> {
objects: ReadableStreamClassObjects<
'js,
ReadableStreamControllerOwned<'js>,
ReadableStreamDefaultReaderOwned<'js>,
>,
prevent_cancel: bool,
is_finished: Rc<AtomicBool>,
ongoing_promise: Option<Promise<'js>>,
promise_primordials: PromisePrimordials<'js>,
end_of_iteration: Symbol<'js>,
}
impl<'js> Trace<'js> for ReadableStreamAsyncIterator<'js> {
fn trace<'a>(&self, tracer: Tracer<'a, 'js>) {
Trace::<'js>::trace(&self.objects, tracer);
if let Some(ongoing_promise) = &self.ongoing_promise {
ongoing_promise.trace(tracer);
}
Trace::<'js>::trace(&self.end_of_iteration, tracer);
}
}
unsafe impl<'js> JsLifetime<'js> for ReadableStreamAsyncIterator<'js> {
type Changed<'to> = ReadableStreamAsyncIterator<'to>;
}
impl<'js> ReadableStreamAsyncIterator<'js> {
pub(super) fn new(
ctx: Ctx<'js>,
objects: ReadableStreamClassObjects<
'js,
ReadableStreamControllerOwned<'js>,
ReadableStreamDefaultReaderOwned<'js>,
>,
promise_primordials: PromisePrimordials<'js>,
prevent_cancel: bool,
) -> Result<Class<'js, Self>> {
let end_of_iteration = IteratorPrimordials::get(&ctx)?.end_of_iteration.clone();
Class::instance(
ctx,
Self {
objects,
prevent_cancel,
is_finished: Rc::new(AtomicBool::new(false)),
ongoing_promise: None,
promise_primordials,
end_of_iteration,
},
)
}
}
impl<'js> JsClass<'js> for ReadableStreamAsyncIterator<'js> {
const NAME: &'static str = "ReadableStreamAsyncIterator";
type Mutable = rquickjs::class::Writable;
fn prototype(ctx: &Ctx<'js>) -> Result<Option<Object<'js>>> {
use rquickjs::class::impl_::MethodImplementor;
let proto = Object::new(ctx.clone())?;
let primordial = IteratorPrimordials::get(ctx)?;
proto.set_prototype(Some(&primordial.async_iterator_prototype))?;
let implementor = rquickjs::class::impl_::MethodImpl::<Self>::new();
implementor.implement(&proto)?;
let next_fn: Function<'js> = proto.get("next")?;
next_fn.set_name("next")?;
let return_fn: Function<'js> = proto.get("return")?;
return_fn.set_name("return")?;
return_fn.set_length(1)?;
let define_property: Function<'js> = ctx
.globals()
.get::<_, Object<'js>>("Object")?
.get("defineProperty")?;
for name in ["next", "return"] {
let value: Value<'js> = proto.get(name)?;
let desc = Object::new(ctx.clone())?;
desc.set("value", value)?;
desc.set("writable", true)?;
desc.set("enumerable", true)?;
desc.set("configurable", true)?;
define_property.call::<_, ()>((proto.clone(), name, desc))?;
}
Ok(Some(proto))
}
fn constructor(ctx: &Ctx<'js>) -> Result<Option<Constructor<'js>>> {
use rquickjs::class::impl_::ConstructorCreator;
let implementor = rquickjs::class::impl_::ConstructorCreate::<Self>::new();
(&implementor).create_constructor(ctx)
}
}
impl<'js> IntoJs<'js> for ReadableStreamAsyncIterator<'js> {
fn into_js(self, ctx: &Ctx<'js>) -> Result<Value<'js>> {
let cls = Class::<Self>::instance(ctx.clone(), self)?;
IntoJs::into_js(cls, ctx)
}
}
impl<'js> FromJs<'js> for ReadableStreamAsyncIterator<'js>
where
for<'a> CloneWrapper<'a, Self>: CloneTrait<Self>,
{
fn from_js(ctx: &Ctx<'js>, value: Value<'js>) -> Result<Self> {
use rquickjs::class::impl_::CloneTrait;
let value = Class::<Self>::from_js(ctx, value)?;
let borrow = value.try_borrow()?;
Ok(CloneWrapper(&*borrow).wrap_clone())
}
}
#[methods]
impl<'js> ReadableStreamAsyncIterator<'js> {
fn next(ctx: Ctx<'js>, iterator: This<OwnedBorrowMut<'js, Self>>) -> Result<Promise<'js>> {
let is_finished = iterator.is_finished.clone();
let next_steps = move |ctx: Ctx<'js>, iterator: &Self, iterator_class: Class<'js, Self>| {
if is_finished.load(Ordering::Acquire) {
return promise_resolved_with(
&ctx,
&iterator.promise_primordials,
Ok(ReadableStreamReadResult {
value: None,
done: true,
}
.into_js(&ctx)?),
);
}
let next_promise = Self::next_steps(&ctx, iterator)?;
upon_promise(
ctx,
next_promise,
move |ctx, result: std::result::Result<Value<'js>, _>| {
let mut iterator = OwnedBorrowMut::from_class(iterator_class);
match result {
Ok(next) => {
iterator.ongoing_promise = None;
if next.as_symbol() == Some(&iterator.end_of_iteration) {
iterator.is_finished.store(true, Ordering::Release);
Ok(ReadableStreamReadResult {
value: None,
done: true,
})
} else {
Ok(ReadableStreamReadResult {
value: Some(next),
done: false,
})
}
},
Err(reason) => {
iterator.ongoing_promise = None;
iterator.is_finished.store(true, Ordering::Release);
Err(ctx.throw(reason))
},
}
},
)
};
let (iterator_class, mut iterator) = class_from_owned_borrow_mut(iterator.0);
let ongoing_promise = iterator.ongoing_promise.take();
let ongoing_promise = match ongoing_promise {
Some(ongoing_promise) => upon_promise(
ctx,
ongoing_promise,
move |ctx, _: std::result::Result<Value<'js>, _>| {
let iterator = OwnedBorrow::from_class(iterator_class.clone());
next_steps(ctx, &iterator, iterator_class)
},
)?,
None => next_steps(ctx, &iterator, iterator_class)?,
};
Ok(iterator.ongoing_promise.insert(ongoing_promise).clone())
}
#[qjs(rename = "return")]
fn r#return(
ctx: Ctx<'js>,
iterator: This<OwnedBorrowMut<'js, Self>>,
value: Opt<Value<'js>>,
) -> Result<Promise<'js>> {
let is_finished = iterator.is_finished.clone();
let value = value.0.unwrap_or_undefined(&ctx);
let return_steps = {
let value = value.clone();
move |ctx: Ctx<'js>, iterator: &Self| {
if is_finished.swap(true, Ordering::AcqRel) {
return promise_resolved_with(
&ctx,
&iterator.promise_primordials,
Ok(ReadableStreamReadResult {
value: Some(value),
done: true,
}
.into_js(&ctx)?),
);
}
Self::return_steps(ctx.clone(), iterator, value)
}
};
let (iterator_class, mut iterator) = class_from_owned_borrow_mut(iterator.0);
let ongoing_promise = iterator.ongoing_promise.take();
let ongoing_promise = match ongoing_promise {
Some(ongoing_promise) => upon_promise(
ctx.clone(),
ongoing_promise,
move |ctx, _: std::result::Result<Value<'js>, _>| {
let iterator = OwnedBorrow::from_class(iterator_class.clone());
return_steps(ctx, &iterator)
},
)?,
None => return_steps(ctx.clone(), &iterator)?,
};
iterator.ongoing_promise = Some(ongoing_promise.clone());
upon_promise_fulfilment(ctx, ongoing_promise, move |_, ()| {
Ok(ReadableStreamReadResult {
value: Some(value),
done: true,
})
})
}
}
impl<'js> ReadableStreamAsyncIterator<'js> {
fn next_steps(ctx: &Ctx<'js>, iterator: &Self) -> Result<Promise<'js>> {
let objects = iterator.objects.clone();
let promise = ResolveablePromise::new(ctx)?;
#[derive(Trace)]
struct ReadRequest<'js> {
promise: ResolveablePromise<'js>,
end_of_iteration: Symbol<'js>,
}
impl<'js> ReadableStreamReadRequest<'js> for ReadRequest<'js> {
fn chunk_steps(
&self,
objects: ReadableStreamDefaultReaderObjects<'js>,
chunk: Value<'js>,
) -> Result<ReadableStreamDefaultReaderObjects<'js>> {
self.promise.resolve(chunk)?;
Ok(objects)
}
fn close_steps(
&self,
_ctx: &Ctx<'js>,
mut objects: ReadableStreamDefaultReaderObjects<'js>,
) -> Result<ReadableStreamDefaultReaderObjects<'js>> {
objects =
ReadableStreamDefaultReader::readable_stream_default_reader_release(objects)?;
self.promise.resolve(self.end_of_iteration.clone())?;
Ok(objects)
}
fn error_steps(
&self,
mut objects: ReadableStreamDefaultReaderObjects<'js>,
reason: Value<'js>,
) -> Result<ReadableStreamDefaultReaderObjects<'js>> {
objects =
ReadableStreamDefaultReader::readable_stream_default_reader_release(objects)?;
self.promise.reject(reason)?;
Ok(objects)
}
}
let objects = ReadableStreamObjects::from_class(objects);
ReadableStreamDefaultReader::readable_stream_default_reader_read(
ctx,
objects,
ReadRequest {
promise: promise.clone(),
end_of_iteration: iterator.end_of_iteration.clone(),
},
)?;
Ok(promise.promise)
}
fn return_steps(ctx: Ctx<'js>, iterator: &Self, arg: Value<'js>) -> Result<Promise<'js>> {
let objects = ReadableStreamObjects::from_class(iterator.objects.clone());
if !iterator.prevent_cancel {
let (result, objects) =
ReadableStreamGenericReader::readable_stream_reader_generic_cancel(
ctx.clone(),
objects,
arg,
)?;
ReadableStreamDefaultReader::readable_stream_default_reader_release(objects)?;
return Ok(result);
}
ReadableStreamDefaultReader::readable_stream_default_reader_release(objects)?;
Ok(iterator
.promise_primordials
.promise_resolved_with_undefined
.clone())
}
}
#[derive(Clone, JsLifetime, Trace)]
pub(crate) struct IteratorPrimordials<'js> {
end_of_iteration: Symbol<'js>,
sync_to_async_iterator: Function<'js>,
async_iterator_prototype: Object<'js>,
}
impl<'js> Primordial<'js> for IteratorPrimordials<'js> {
fn new(ctx: &Ctx<'js>) -> Result<Self>
where
Self: Sized,
{
let sync_to_async_iterator = ctx.eval::<Function<'js>, _>(
r#"
(syncIterable) => (async function* () {
return yield* syncIterable;
})()
"#,
)?;
let async_iterator_prototype = ctx
.eval::<Object<'js>, _>("(async function* () {})()")?
.get_prototype()
.as_ref()
.and_then(Object::get_prototype)
.as_ref()
.and_then(Object::get_prototype)
.expect("async iterator prototype not found");
Ok(Self {
end_of_iteration: Symbol::new_global(ctx.clone(), "async iterator end of iteration")?,
sync_to_async_iterator,
async_iterator_prototype,
})
}
}
fn get_method<'js>(
ctx: &Ctx<'js>,
value: Value<'js>,
property: impl IntoAtom<'js>,
) -> Result<Option<Function<'js>>> {
let func = get_v(ctx, value, property)?;
if func.is_undefined() || func.is_null() {
return Ok(None);
}
match func.into_function() {
None => Err(Exception::throw_type(ctx, "not a function")),
Some(func) => Ok(Some(func)),
}
}
fn get_v<'js>(
ctx: &Ctx<'js>,
value: Value<'js>,
property: impl IntoAtom<'js>,
) -> Result<Value<'js>> {
let o: Object<'js> = to_object(ctx, value)?;
o.get(property)
}
fn to_object<'js>(ctx: &Ctx<'js>, value: Value<'js>) -> Result<Object<'js>> {
let base_primordials = BasePrimordials::get(ctx)?;
match value.type_of() {
Type::Bool => base_primordials.constructor_bool.construct((value,))?,
Type::Int | Type::Float => base_primordials.constructor_number.construct((value,))?,
Type::String => base_primordials.constructor_string.construct((value,))?,
Type::Symbol => base_primordials.constructor_object.call((value,))?,
Type::BigInt => base_primordials.constructor_object.call((value,))?,
typ if typ.interpretable_as(Type::Object) => Ok(value.into_object().unwrap()),
typ => Err(Exception::throw_type(
ctx,
&format!("{typ} cannot be converted to an object"),
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test::test_sync_with;
use rquickjs::BigInt;
#[tokio::test]
async fn test_to_object() {
test_sync_with(|ctx| {
BasePrimordials::init(&ctx)?;
let good_values: [Value; 7] = [
Value::new_bool(ctx.clone(), false),
Value::new_int(ctx.clone(), 123),
Value::new_float(ctx.clone(), 1.5),
rquickjs::String::from_str(ctx.clone(), "abc")?.into_value(),
Symbol::new_global(ctx.clone(), "def")?.into_value(),
BigInt::from_i64(ctx.clone(), 123456)?.into_value(),
Object::new(ctx.clone())?.into_value(),
];
for value in good_values {
to_object(&ctx, value)?;
}
let bad_values: [Value; 3] = [
Value::new_uninitialized(ctx.clone()),
Value::new_undefined(ctx.clone()),
Value::new_null(ctx.clone()),
];
for value in bad_values {
let ty = value.type_of();
if to_object(&ctx, value).is_ok() {
panic!("Values of type {ty} should not be convertible to object")
}
}
Ok(())
})
.await;
}
}