use crate::utils::option::{Null, Undefined};
use rquickjs::{
class::{OwnedBorrow, OwnedBorrowMut, Trace},
methods,
prelude::{Opt, This},
Class, Ctx, Error, Exception, JsLifetime, Object, Promise, Result, Value,
};
use std::{future, pin::Pin, rc::Rc};
pub enum NativePullResult<'js> {
Ready(Value<'js>),
Eof,
Pending(Pin<Box<dyn future::Future<Output = Result<Option<Value<'js>>>> + 'js>>),
}
pub type NativePullFn<'js> = dyn Fn(&Ctx<'js>) -> Result<NativePullResult<'js>> + 'js;
pub struct NativePull<'js>(pub Rc<NativePullFn<'js>>);
impl<'js> Clone for NativePull<'js> {
fn clone(&self) -> Self {
Self(self.0.clone())
}
}
unsafe impl<'js> JsLifetime<'js> for NativePull<'js> {
type Changed<'to> = NativePull<'to>;
}
impl<'js> Trace<'js> for NativePull<'js> {
fn trace<'a>(&self, _: rquickjs::class::Tracer<'a, 'js>) {}
}
use crate::stream_web::{
queuing_strategy::{SizeAlgorithm, SizeValue},
readable::{
byte_controller::ReadableByteStreamControllerOwned,
controller::{
ReadableStreamController, ReadableStreamControllerClass, ReadableStreamControllerOwned,
},
default_reader::{ReadableStreamDefaultReaderOrUndefined, ReadableStreamReadRequest},
objects::{
ReadableStreamClassObjects, ReadableStreamDefaultControllerObjects,
ReadableStreamDefaultReaderObjects, ReadableStreamObjects,
},
reader::ReadableStreamReader,
stream::{
algorithms::{CancelAlgorithm, PullAlgorithm, StartAlgorithm},
source::UnderlyingSource,
ReadableStream, ReadableStreamClass, ReadableStreamOwned, ReadableStreamState,
},
},
utils::{
class_from_owned_borrow_mut,
promise::{promise_resolved_with, upon_promise},
queue::QueueWithSizes,
UnwrapOrUndefined,
},
};
#[derive(JsLifetime, Trace)]
#[rquickjs::class]
pub struct ReadableStreamDefaultController<'js> {
cancel_algorithm: Option<CancelAlgorithm<'js>>,
pub(super) close_requested: bool,
pull_again: bool,
pull_algorithm: Option<PullAlgorithm<'js>>,
pub(crate) pulling: bool,
pub(crate) container: QueueWithSizes<'js>,
started: bool,
strategy_hwm: f64,
strategy_size_algorithm: Option<SizeAlgorithm<'js>>,
pub(super) stream: ReadableStreamClass<'js>,
pub native_pull: Option<NativePull<'js>>,
#[qjs(skip_trace)]
pub(super) is_owning_type: bool,
}
impl<'js> Drop for ReadableStreamDefaultController<'js> {
fn drop(&mut self) {
self.native_pull = None;
}
}
pub type ReadableStreamDefaultControllerClass<'js> =
Class<'js, ReadableStreamDefaultController<'js>>;
pub(super) type ReadableStreamDefaultControllerOwned<'js> =
OwnedBorrowMut<'js, ReadableStreamDefaultController<'js>>;
impl<'js> ReadableStreamDefaultController<'js> {
pub(super) fn set_up_readable_stream_default_controller_from_underlying_source(
ctx: Ctx<'js>,
stream: ReadableStreamOwned<'js>,
underlying_source: Null<Undefined<Object<'js>>>,
underlying_source_dict: UnderlyingSource<'js>,
high_water_mark: f64,
size_algorithm: SizeAlgorithm<'js>,
is_owning_type: bool,
) -> Result<()> {
let (start_algorithm, pull_algorithm, cancel_algorithm) = (
underlying_source_dict
.start
.map(|f| StartAlgorithm::Function {
f,
underlying_source: underlying_source.clone(),
})
.unwrap_or(StartAlgorithm::ReturnUndefined),
underlying_source_dict
.pull
.map(|f| PullAlgorithm::Function {
f,
underlying_source: underlying_source.clone(),
})
.unwrap_or(PullAlgorithm::ReturnPromiseUndefined),
underlying_source_dict
.cancel
.map(|f| CancelAlgorithm::Function {
f,
underlying_source,
})
.unwrap_or(CancelAlgorithm::ReturnPromiseUndefined),
);
Self::set_up_readable_stream_default_controller(
ctx.clone(),
stream,
start_algorithm,
pull_algorithm,
cancel_algorithm,
high_water_mark,
size_algorithm,
is_owning_type,
)?;
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(super) fn set_up_readable_stream_default_controller(
ctx: Ctx<'js>,
stream: ReadableStreamOwned<'js>,
start_algorithm: StartAlgorithm<'js>,
pull_algorithm: PullAlgorithm<'js>,
cancel_algorithm: CancelAlgorithm<'js>,
high_water_mark: f64,
size_algorithm: SizeAlgorithm<'js>,
is_owning_type: bool,
) -> Result<Class<'js, Self>> {
let (stream_class, mut stream) = class_from_owned_borrow_mut(stream);
let controller = ReadableStreamDefaultController {
stream: stream_class.clone(),
container: QueueWithSizes::new(),
started: false,
close_requested: false,
pull_again: false,
pulling: false,
strategy_size_algorithm: Some(size_algorithm),
strategy_hwm: high_water_mark,
pull_algorithm: Some(pull_algorithm),
cancel_algorithm: Some(cancel_algorithm),
native_pull: None,
is_owning_type,
};
let controller_class = Class::instance(ctx.clone(), controller)?;
stream.controller = ReadableStreamControllerClass::ReadableStreamDefaultController(
controller_class.clone(),
);
let objects = ReadableStreamObjects::new_default(
stream,
OwnedBorrowMut::from_class(controller_class),
);
let promise_primordials = objects.stream.promise_primordials.clone();
let (start_result, objects_class) =
Self::start_algorithm(ctx.clone(), objects, start_algorithm)?;
let start_promise = promise_resolved_with(&ctx, &promise_primordials, Ok(start_result))?;
let _ = upon_promise::<Value<'js>, _>(ctx.clone(), start_promise, {
let objects_class = objects_class.clone();
move |ctx, result| {
let mut objects =
ReadableStreamObjects::from_class_no_reader(objects_class).refresh_reader();
match result {
Ok(_) => {
objects.controller.started = true;
Self::readable_stream_default_controller_call_pull_if_needed(ctx, objects)?;
},
Err(r) => {
Self::readable_stream_default_controller_error(objects, r)?;
},
}
Ok(())
}
})?;
Ok(objects_class.controller)
}
fn readable_stream_default_controller_call_pull_if_needed<
R: ReadableStreamDefaultReaderOrUndefined<'js>,
>(
ctx: Ctx<'js>,
objects: ReadableStreamDefaultControllerObjects<'js, R>,
) -> Result<ReadableStreamDefaultControllerObjects<'js, R>> {
let (should_pull, mut objects) =
ReadableStreamDefaultController::readable_stream_default_controller_should_call_pull(
objects,
);
if !should_pull {
return Ok(objects);
}
if objects.controller.pulling {
objects.controller.pull_again = true;
return Ok(objects);
}
objects.controller.pulling = true;
let (pull_promise, objects_class) = Self::pull_algorithm(ctx.clone(), objects)?;
upon_promise::<Value<'js>, _>(ctx.clone(), pull_promise, {
let objects_class = objects_class.clone();
move |ctx, result| {
let mut objects =
ReadableStreamObjects::from_class_no_reader(objects_class).refresh_reader();
match result {
Ok(_) => {
objects.controller.pulling = false;
if objects.controller.pull_again {
objects.controller.pull_again = false;
Self::readable_stream_default_controller_call_pull_if_needed(
ctx, objects,
)?;
};
Ok(())
},
Err(e) => {
Self::readable_stream_default_controller_error(objects, e)?;
Ok(())
},
}
}
})?;
Ok(ReadableStreamObjects::from_class(objects_class))
}
pub(super) fn readable_stream_default_controller_error<R: ReadableStreamReader<'js>>(
mut objects: ReadableStreamDefaultControllerObjects<'js, R>,
e: Value<'js>,
) -> Result<ReadableStreamDefaultControllerObjects<'js, R>> {
if !matches!(objects.stream.state, ReadableStreamState::Readable) {
return Ok(objects);
};
objects.controller.container.reset_queue();
objects
.controller
.readable_stream_default_controller_clear_algorithms();
ReadableStream::readable_stream_error(objects, e)
}
fn readable_stream_default_controller_should_call_pull<
R: ReadableStreamDefaultReaderOrUndefined<'js>,
>(
mut objects: ReadableStreamDefaultControllerObjects<'js, R>,
) -> (bool, ReadableStreamDefaultControllerObjects<'js, R>) {
if !objects
.controller
.readable_stream_default_controller_can_close_or_enqueue(&objects.stream)
{
return (false, objects);
}
if !objects.controller.started {
return (false, objects);
}
{
let mut ret = false;
objects = objects
.with_some_reader(
|objects| {
if ReadableStream::readable_stream_get_num_read_requests(&objects.reader)
> 0
{
ret = true
}
Ok(objects)
},
Ok,
)
.unwrap();
if ret {
return (true, objects);
}
}
let desired_size = objects.controller
.readable_stream_default_controller_get_desired_size(&objects.stream)
.0
.expect(
"desiredSize should not be null during ReadableStreamDefaultControllerShouldCallPull",
);
if desired_size > 0.0 {
return (true, objects);
}
(false, objects)
}
fn readable_stream_default_controller_clear_algorithms(&mut self) {
self.pull_algorithm = None;
self.cancel_algorithm = None;
self.strategy_size_algorithm = None;
self.native_pull = None;
}
fn readable_stream_default_controller_can_close_or_enqueue(
&self,
stream: &ReadableStream<'js>,
) -> bool {
match stream.state {
ReadableStreamState::Readable if !self.close_requested => true,
_ => false,
}
}
pub(crate) fn readable_stream_default_controller_get_desired_size(
&self,
stream: &ReadableStream<'js>,
) -> Null<f64> {
match stream.state {
ReadableStreamState::Errored(_) => Null(None),
ReadableStreamState::Closed => Null(Some(0.0)),
ReadableStreamState::Readable => {
Null(Some(self.strategy_hwm - self.container.queue_total_size))
},
}
}
pub(super) fn readable_stream_default_controller_close<R: ReadableStreamReader<'js>>(
ctx: Ctx<'js>,
mut objects: ReadableStreamDefaultControllerObjects<'js, R>,
) -> Result<ReadableStreamDefaultControllerObjects<'js, R>> {
if !objects
.controller
.readable_stream_default_controller_can_close_or_enqueue(&objects.stream)
{
return Ok(objects);
}
objects.controller.close_requested = true;
if objects.controller.container.queue.is_empty() {
objects
.controller
.readable_stream_default_controller_clear_algorithms();
objects = ReadableStream::readable_stream_close(ctx, objects)?;
}
Ok(objects)
}
pub(super) fn readable_stream_default_controller_enqueue<
R: ReadableStreamDefaultReaderOrUndefined<'js>,
>(
ctx: Ctx<'js>,
mut objects: ReadableStreamDefaultControllerObjects<'js, R>,
chunk: Value<'js>,
) -> Result<ReadableStreamDefaultControllerObjects<'js, R>> {
if !objects
.controller
.readable_stream_default_controller_can_close_or_enqueue(&objects.stream)
{
return Ok(objects);
}
let mut els = true;
objects = objects.with_some_reader(
|objects| {
if ReadableStream::readable_stream_get_num_read_requests(&objects.reader) > 0 {
els = false;
ReadableStream::readable_stream_fulfill_read_request(
&ctx,
objects,
chunk.clone(),
false,
)
} else {
Ok(objects)
}
},
Ok,
)?;
if els {
let (result, objects_class) =
Self::strategy_size_algorithm(ctx.clone(), objects, chunk.clone());
objects = ReadableStreamObjects::from_class(objects_class);
match result {
Err(Error::Exception) => {
let err = ctx.catch();
Self::readable_stream_default_controller_error(objects, err.clone())?;
return Err(ctx.throw(err));
},
Ok(chunk_size) => {
let enqueue_result = objects
.controller
.container
.enqueue_value_with_size(&ctx, chunk, chunk_size);
match enqueue_result {
Err(Error::Exception) => {
let err = ctx.catch();
Self::readable_stream_default_controller_error(objects, err.clone())?;
return Err(ctx.throw(err));
},
Err(err) => return Err(err),
Ok(()) => {},
}
},
Err(err) => return Err(err),
}
}
Self::readable_stream_default_controller_call_pull_if_needed(ctx, objects)
}
fn start_algorithm<R: ReadableStreamReader<'js>>(
ctx: Ctx<'js>,
objects: ReadableStreamDefaultControllerObjects<'js, R>,
start_algorithm: StartAlgorithm<'js>,
) -> Result<(
Value<'js>,
ReadableStreamClassObjects<'js, OwnedBorrowMut<'js, Self>, R>,
)> {
let objects_class = objects.into_inner();
Ok((
start_algorithm.call(
ctx,
ReadableStreamControllerClass::ReadableStreamDefaultController(
objects_class.controller.clone(),
),
)?,
objects_class,
))
}
fn pull_algorithm<R: ReadableStreamReader<'js>>(
ctx: Ctx<'js>,
objects: ReadableStreamDefaultControllerObjects<'js, R>,
) -> Result<(
Promise<'js>,
ReadableStreamClassObjects<'js, OwnedBorrowMut<'js, Self>, R>,
)> {
let pull_algorithm = objects
.controller
.pull_algorithm
.clone()
.expect("pull algorithm used after ReadableStreamDefaultControllerClearAlgorithms");
let promise_primordials = objects.stream.promise_primordials.clone();
let objects_class = objects.into_inner();
Ok((
pull_algorithm.call(
ctx,
&promise_primordials,
ReadableStreamControllerClass::ReadableStreamDefaultController(
objects_class.controller.clone(),
),
)?,
objects_class,
))
}
fn strategy_size_algorithm<R: ReadableStreamReader<'js>>(
ctx: Ctx<'js>,
objects: ReadableStreamDefaultControllerObjects<'js, R>,
chunk: Value<'js>,
) -> (
Result<SizeValue<'js>>,
ReadableStreamClassObjects<'js, OwnedBorrowMut<'js, Self>, R>,
) {
let strategy_size_algorithm = objects
.controller
.strategy_size_algorithm
.clone()
.expect("size algorithm used after ReadableStreamDefaultControllerClearAlgorithms");
let objects_class = objects.into_inner();
(strategy_size_algorithm.call(ctx, chunk), objects_class)
}
pub(super) fn cancel_algorithm<R: ReadableStreamReader<'js>>(
ctx: Ctx<'js>,
objects: ReadableStreamDefaultControllerObjects<'js, R>,
reason: Value<'js>,
) -> Result<(
Promise<'js>,
ReadableStreamClassObjects<'js, OwnedBorrowMut<'js, Self>, R>,
)> {
let cancel_algorithm =
objects.controller.cancel_algorithm.clone().expect(
"cancel algorithm used after ReadableStreamDefaultControllerClearAlgorithms",
);
let promise_primordials = objects.stream.promise_primordials.clone();
let objects_class = objects.into_inner();
Ok((
cancel_algorithm.call(ctx, &promise_primordials, reason)?,
objects_class,
))
}
}
#[methods(rename_all = "camelCase")]
impl<'js> ReadableStreamDefaultController<'js> {
fn constructor() -> Self {
unimplemented!()
}
#[qjs(constructor)]
fn new(ctx: Ctx<'js>) -> Result<Class<'js, Self>> {
Err(Exception::throw_type(&ctx, "Illegal constructor"))
}
#[qjs(get)]
fn desired_size(&self) -> Null<f64> {
let stream = OwnedBorrow::from_class(self.stream.clone());
self.readable_stream_default_controller_get_desired_size(&stream)
}
fn close(ctx: Ctx<'js>, controller: This<OwnedBorrowMut<'js, Self>>) -> Result<()> {
let objects = ReadableStreamObjects::from_default_controller(controller.0);
if !objects
.controller
.readable_stream_default_controller_can_close_or_enqueue(&objects.stream)
{
return Err(Exception::throw_type(
&ctx,
"The stream is not in a state that permits close",
));
}
Self::readable_stream_default_controller_close(ctx, objects)?;
Ok(())
}
fn enqueue(
ctx: Ctx<'js>,
controller: This<OwnedBorrowMut<'js, Self>>,
chunk: Opt<Value<'js>>,
options: Opt<Value<'js>>,
) -> Result<()> {
let mut transfer_list: Option<rquickjs::Array<'js>> = None;
if let Some(opts) = options.0.as_ref().and_then(|v| v.as_object()) {
transfer_list = opts.get::<_, Option<rquickjs::Array<'js>>>("transfer")?;
}
let has_transfer_items = transfer_list.as_ref().is_some_and(|arr| !arr.is_empty());
if has_transfer_items && !controller.is_owning_type {
return Err(Exception::throw_type(&ctx, "transfer list is not empty"));
}
let chunk_value = chunk.0.clone().unwrap_or_undefined(&ctx);
let transferred_chunk = if has_transfer_items && controller.is_owning_type {
transfer_owning_chunk(&ctx, chunk_value.clone(), &transfer_list.unwrap())?
} else {
chunk_value
};
let objects = ReadableStreamObjects::from_default_controller(controller.0);
if !objects
.controller
.readable_stream_default_controller_can_close_or_enqueue(&objects.stream)
{
return Err(Exception::throw_type(
&ctx,
"The stream is not in a state that permits enqueue",
));
}
objects.with_reader(
|objects| {
Self::readable_stream_default_controller_enqueue(
ctx.clone(),
objects,
transferred_chunk.clone(),
)
},
|_| panic!("Default controller must not have byob reader"),
|objects| {
Self::readable_stream_default_controller_enqueue(
ctx.clone(),
objects,
transferred_chunk.clone(),
)
},
)?;
Ok(())
}
fn error(
ctx: Ctx<'js>,
controller: This<OwnedBorrowMut<'js, Self>>,
e: Opt<Value<'js>>,
) -> Result<()> {
let objects = ReadableStreamObjects::from_default_controller(controller.0);
Self::readable_stream_default_controller_error(objects, e.0.unwrap_or_undefined(&ctx))?;
Ok(())
}
}
impl<'js> ReadableStreamController<'js> for ReadableStreamDefaultControllerOwned<'js> {
type Class = ReadableStreamDefaultControllerClass<'js>;
fn with_controller<C, O>(
self,
ctx: C,
default: impl FnOnce(
C,
ReadableStreamDefaultControllerOwned<'js>,
) -> Result<(O, ReadableStreamDefaultControllerOwned<'js>)>,
_: impl FnOnce(
C,
ReadableByteStreamControllerOwned<'js>,
) -> Result<(O, ReadableByteStreamControllerOwned<'js>)>,
) -> Result<(O, Self)> {
let (ctx, reader) = default(ctx, self)?;
Ok((ctx, reader))
}
fn into_inner(self) -> Self::Class {
OwnedBorrowMut::into_inner(self)
}
fn from_class(class: Self::Class) -> Self {
OwnedBorrowMut::from_class(class)
}
fn into_erased(self) -> ReadableStreamControllerOwned<'js> {
ReadableStreamControllerOwned::ReadableStreamDefaultController(self)
}
fn try_from_erased(erased: ReadableStreamControllerOwned<'js>) -> Option<Self> {
match erased {
ReadableStreamControllerOwned::ReadableStreamDefaultController(r) => Some(r),
ReadableStreamControllerOwned::ReadableStreamByteController(_) => None,
}
}
fn pull_steps(
ctx: &Ctx<'js>,
mut objects: ReadableStreamDefaultReaderObjects<'js, Self>,
read_request: impl ReadableStreamReadRequest<'js> + 'js,
) -> Result<ReadableStreamDefaultReaderObjects<'js, Self>> {
if !objects.controller.container.queue.is_empty() {
let chunk = objects.controller.container.dequeue_value();
if objects.controller.close_requested && objects.controller.container.queue.is_empty() {
objects
.controller
.readable_stream_default_controller_clear_algorithms();
objects = ReadableStream::readable_stream_close(ctx.clone(), objects)?;
} else {
objects =
ReadableStreamDefaultController::readable_stream_default_controller_call_pull_if_needed(
ctx.clone(),
objects,
)?;
}
read_request.chunk_steps_typed(objects, chunk)
} else {
objects
.stream
.readable_stream_add_read_request(&mut objects.reader, read_request);
ReadableStreamDefaultController::readable_stream_default_controller_call_pull_if_needed(
ctx.clone(),
objects,
)
}
}
fn cancel_steps<R: ReadableStreamReader<'js>>(
ctx: &Ctx<'js>,
mut objects: ReadableStreamObjects<'js, Self, R>,
reason: Value<'js>,
) -> Result<(Promise<'js>, ReadableStreamObjects<'js, Self, R>)> {
objects.controller.container.reset_queue();
let (result, objects_class) =
ReadableStreamDefaultController::cancel_algorithm(ctx.clone(), objects, reason)?;
objects = ReadableStreamObjects::from_class(objects_class);
objects
.controller
.readable_stream_default_controller_clear_algorithms();
Ok((result, objects))
}
fn release_steps(&mut self) {}
}
pub fn readable_stream_default_controller_enqueue_value<'js>(
ctx: Ctx<'js>,
controller: ReadableStreamDefaultControllerClass<'js>,
chunk: Value<'js>,
) -> Result<()> {
let objects =
ReadableStreamObjects::from_default_controller(OwnedBorrowMut::from_class(controller));
if !objects
.controller
.readable_stream_default_controller_can_close_or_enqueue(&objects.stream)
{
return Ok(()); }
objects.with_reader(
|objects| {
ReadableStreamDefaultController::readable_stream_default_controller_enqueue(
ctx.clone(),
objects,
chunk.clone(),
)
},
|_| panic!("Default controller must not have byob reader"),
|objects| {
ReadableStreamDefaultController::readable_stream_default_controller_enqueue(
ctx.clone(),
objects,
chunk.clone(),
)
},
)?;
Ok(())
}
pub fn readable_stream_default_controller_close_stream<'js>(
ctx: Ctx<'js>,
controller: ReadableStreamDefaultControllerClass<'js>,
) -> Result<()> {
let objects =
ReadableStreamObjects::from_default_controller(OwnedBorrowMut::from_class(controller));
if !objects
.controller
.readable_stream_default_controller_can_close_or_enqueue(&objects.stream)
{
return Ok(());
}
ReadableStreamDefaultController::readable_stream_default_controller_close(ctx, objects)?;
Ok(())
}
pub fn readable_stream_default_controller_error_stream<'js>(
controller: ReadableStreamDefaultControllerClass<'js>,
error: Value<'js>,
) -> Result<()> {
let objects =
ReadableStreamObjects::from_default_controller(OwnedBorrowMut::from_class(controller));
objects.with_reader(
|objects| {
ReadableStreamDefaultController::readable_stream_default_controller_error(
objects,
error.clone(),
)
},
|_| panic!("Default controller must not have byob reader"),
|objects| {
ReadableStreamDefaultController::readable_stream_default_controller_error(
objects,
error.clone(),
)
},
)?;
Ok(())
}
fn transfer_owning_chunk<'js>(
ctx: &Ctx<'js>,
chunk: Value<'js>,
transfer_list: &rquickjs::Array<'js>,
) -> Result<Value<'js>> {
use rquickjs::ArrayBuffer;
let mut chunk_replacement: Option<Value<'js>> = None;
for v in transfer_list.iter::<Value<'js>>() {
let v = v?;
let Some(ab) = ArrayBuffer::from_value(v.clone()) else {
return Err(rquickjs::Exception::throw_type(
ctx,
"transfer list item is not an ArrayBuffer",
));
};
let is_chunk = chunk == v;
let transfer_fn: rquickjs::Function<'js> = ab.as_object().get("transfer")?;
let new_buf: Value<'js> = transfer_fn.call((rquickjs::function::This(ab.clone()),))?;
if is_chunk && chunk_replacement.is_none() {
chunk_replacement = Some(new_buf);
}
}
Ok(chunk_replacement.unwrap_or(chunk))
}