use crate::abort::{AbortController, AbortSignal};
use crate::utils::{
option::{Null, Undefined},
primordials::Primordial,
};
use rquickjs::{
class::{JsClass, OwnedBorrowMut, Trace},
function::Constructor,
methods,
prelude::{Opt, This},
Class, Ctx, Error, Exception, Function, JsLifetime, Object, Promise, Result, Symbol, Value,
};
use crate::stream_web::{
queuing_strategy::{SizeAlgorithm, SizeValue},
transform::controller::TransformStreamDefaultControllerClass,
transform::stream::TransformStreamClass,
utils::{
class_from_owned_borrow_mut,
promise::{promise_resolved_with, upon_promise, PromisePrimordials},
queue::QueueWithSizes,
UnwrapOrUndefined,
},
writable::{
default_writer::WritableStreamDefaultWriterOwned,
objects::{WritableStreamClassObjects, WritableStreamObjects},
stream::{
sink::UnderlyingSink, WritableStream, WritableStreamClass, WritableStreamOwned,
WritableStreamState,
},
writer::{UndefinedWriter, WritableStreamWriter},
},
};
#[rquickjs::class]
#[derive(JsLifetime, Trace)]
pub(crate) struct WritableStreamDefaultController<'js> {
abort_algorithm: Option<WritableAbortAlgorithm<'js>>,
close_algorithm: Option<WritableCloseAlgorithm<'js>>,
container: QueueWithSizes<'js>,
pub(super) started: bool,
strategy_hwm: f64,
strategy_size_algorithm: Option<SizeAlgorithm<'js>>,
pub(super) abort_controller: Class<'js, AbortController<'js>>,
pub(super) stream: WritableStreamClass<'js>,
write_algorithm: Option<WritableWriteAlgorithm<'js>>,
primordials: WritableStreamDefaultControllerPrimordials<'js>,
}
pub(crate) type WritableStreamDefaultControllerClass<'js> =
Class<'js, WritableStreamDefaultController<'js>>;
pub(crate) type WritableStreamDefaultControllerOwned<'js> =
OwnedBorrowMut<'js, WritableStreamDefaultController<'js>>;
impl<'js> WritableStreamDefaultController<'js> {
pub(super) fn set_up_writable_stream_default_controller_from_underlying_sink(
ctx: Ctx<'js>,
stream: WritableStreamOwned<'js>,
underlying_sink: Null<Undefined<Object<'js>>>,
underlying_sink_dict: UnderlyingSink<'js>,
high_water_mark: f64,
size_algorithm: SizeAlgorithm<'js>,
) -> Result<()> {
let (start_algorithm, write_algorithm, close_algorithm, abort_algorithm) = (
underlying_sink_dict
.start
.map(|f| WritableStartAlgorithm::Function {
f,
underlying_sink: underlying_sink.clone(),
})
.unwrap_or(WritableStartAlgorithm::ReturnUndefined),
underlying_sink_dict
.write
.map(|f| WritableWriteAlgorithm::Function {
f,
underlying_sink: underlying_sink.clone(),
})
.unwrap_or(WritableWriteAlgorithm::ReturnPromiseUndefined),
underlying_sink_dict
.close
.map(|f| WritableCloseAlgorithm::Function {
f,
underlying_sink: underlying_sink.clone(),
})
.unwrap_or(WritableCloseAlgorithm::ReturnPromiseUndefined),
underlying_sink_dict
.abort
.map(|f| WritableAbortAlgorithm::Function {
f,
underlying_sink: underlying_sink.clone(),
})
.unwrap_or(WritableAbortAlgorithm::ReturnPromiseUndefined),
);
Self::set_up_writable_stream_default_controller(
ctx,
stream,
start_algorithm,
write_algorithm,
close_algorithm,
abort_algorithm,
high_water_mark,
size_algorithm,
)
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn set_up_writable_stream_default_controller(
ctx: Ctx<'js>,
stream: WritableStreamOwned<'js>,
start_algorithm: WritableStartAlgorithm<'js>,
write_algorithm: WritableWriteAlgorithm<'js>,
close_algorithm: WritableCloseAlgorithm<'js>,
abort_algorithm: WritableAbortAlgorithm<'js>,
high_water_mark: f64,
size_algorithm: SizeAlgorithm<'js>,
) -> Result<()> {
let (stream_class, mut stream) = class_from_owned_borrow_mut(stream);
let controller = Self {
stream: stream_class,
container: QueueWithSizes::new(),
abort_controller: Class::instance(ctx.clone(), AbortController::new(ctx.clone())?)?,
started: false,
strategy_size_algorithm: Some(size_algorithm),
strategy_hwm: high_water_mark,
write_algorithm: Some(write_algorithm),
close_algorithm: Some(close_algorithm),
abort_algorithm: Some(abort_algorithm),
primordials: WritableStreamDefaultControllerPrimordials::get(&ctx)?.clone(),
};
let controller_class = Class::instance(ctx.clone(), controller)?;
stream.controller = Some(controller_class.clone());
let objects = WritableStreamObjects::from_stream(stream);
let backpressure = objects
.controller
.writable_stream_default_controller_get_backpressure();
let objects = WritableStream::writable_stream_update_backpressure(
ctx.clone(),
objects,
backpressure,
)?;
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, {
move |ctx, result| {
let mut objects =
WritableStreamObjects::from_class_no_writer(objects_class).refresh_writer();
match result {
Ok(_) => {
objects.controller.started = true;
Self::writable_stream_default_controller_advance_queue_if_needed(
ctx, objects,
)?;
},
Err(r) => {
objects.controller.started = true;
WritableStream::writable_stream_deal_with_rejection(ctx, objects, r)?;
},
}
Ok(())
}
})?;
Ok(())
}
pub(super) fn writable_stream_default_controller_close<W: WritableStreamWriter<'js>>(
ctx: Ctx<'js>,
mut objects: WritableStreamObjects<'js, W>,
) -> Result<WritableStreamObjects<'js, W>> {
let close_sentinel = objects
.controller
.primordials
.close_sentinel
.as_value()
.clone();
objects.controller.container.enqueue_value_with_size(
&ctx,
close_sentinel,
SizeValue::Native(0.0),
)?;
objects = Self::writable_stream_default_controller_advance_queue_if_needed(ctx, objects)?;
Ok(objects)
}
pub(super) fn writable_stream_default_controller_get_desired_size(&self) -> f64 {
self.strategy_hwm - self.container.queue_total_size
}
pub fn writable_stream_default_controller_get_backpressure(&self) -> bool {
let desired_size = self.writable_stream_default_controller_get_desired_size();
desired_size <= 0.0
}
pub(super) fn writable_stream_default_controller_get_chunk_size(
ctx: Ctx<'js>,
mut objects: WritableStreamObjects<'js, WritableStreamDefaultWriterOwned<'js>>,
chunk: Value<'js>,
) -> Result<(
SizeValue<'js>,
WritableStreamObjects<'js, WritableStreamDefaultWriterOwned<'js>>,
)> {
let (return_value, objects_class) =
Self::strategy_size_algorithm(ctx.clone(), objects, chunk);
match return_value {
Ok(chunk_size) => {
objects = WritableStreamObjects::from_class(objects_class);
Ok((chunk_size, objects))
},
Err(Error::Exception) => {
let reason = ctx.catch();
objects = WritableStreamObjects::from_class(objects_class);
objects = Self::writable_stream_default_controller_error_if_needed(
ctx.clone(),
objects,
reason,
)?;
Ok((SizeValue::Native(1.0), objects))
},
Err(err) => Err(err),
}
}
fn writable_stream_default_controller_error_if_needed(
ctx: Ctx<'js>,
objects: WritableStreamObjects<'js, WritableStreamDefaultWriterOwned<'js>>,
error: Value<'js>,
) -> Result<WritableStreamObjects<'js, WritableStreamDefaultWriterOwned<'js>>> {
if let WritableStreamState::Writable = objects.stream.state {
Self::writable_stream_default_controller_error(ctx, objects, error)
} else {
Ok(objects)
}
}
pub(crate) fn writable_stream_default_controller_error<W: WritableStreamWriter<'js>>(
ctx: Ctx<'js>,
mut objects: WritableStreamObjects<'js, W>,
reason: Value<'js>,
) -> Result<WritableStreamObjects<'js, W>> {
objects
.controller
.writable_stream_default_controller_clear_algorithms();
objects = WritableStream::writable_stream_start_erroring(ctx, objects, reason)?;
Ok(objects)
}
fn writable_stream_default_controller_clear_algorithms(&mut self) {
self.write_algorithm = None;
self.close_algorithm = None;
self.abort_algorithm = None;
self.strategy_size_algorithm = None;
}
pub(super) fn writable_stream_default_controller_write(
ctx: Ctx<'js>,
mut objects: WritableStreamObjects<'js, WritableStreamDefaultWriterOwned<'js>>,
chunk: Value<'js>,
chunk_size: SizeValue<'js>,
) -> Result<WritableStreamObjects<'js, WritableStreamDefaultWriterOwned<'js>>> {
let enqueue_result = objects
.controller
.container
.enqueue_value_with_size(&ctx, chunk, chunk_size);
match enqueue_result {
Err(Error::Exception) => {
let reason = ctx.catch();
objects =
Self::writable_stream_default_controller_error_if_needed(ctx, objects, reason)?;
return Ok(objects);
},
Err(err) => return Err(err),
Ok(()) => {},
}
if !objects.stream.writable_stream_close_queued_or_in_flight()
&& matches!(objects.stream.state, WritableStreamState::Writable)
{
let backpressure = objects
.controller
.writable_stream_default_controller_get_backpressure();
objects = WritableStream::writable_stream_update_backpressure(
ctx.clone(),
objects,
backpressure,
)?;
}
let objects =
Self::writable_stream_default_controller_advance_queue_if_needed(ctx, objects)?;
Ok(objects)
}
fn writable_stream_default_controller_advance_queue_if_needed<W: WritableStreamWriter<'js>>(
ctx: Ctx<'js>,
objects: WritableStreamObjects<'js, W>,
) -> Result<WritableStreamObjects<'js, W>> {
if !objects.controller.started || objects.stream.in_flight_write_request.is_some() {
return Ok(objects);
}
if let WritableStreamState::Erroring(ref stored_error) = objects.stream.state {
let stored_error = stored_error.clone();
return WritableStream::writable_stream_finish_erroring(ctx, objects, stored_error);
}
let value = match objects.controller.container.queue.front() {
None => {
return Ok(objects);
},
Some(value) => value.clone(),
};
if value.value.as_symbol() == Some(&objects.controller.primordials.close_sentinel) {
Self::writable_stream_default_controller_process_close(ctx, objects)
} else {
Self::writable_stream_default_controller_process_write(ctx, objects, value.value)
}
}
fn writable_stream_default_controller_process_close<W: WritableStreamWriter<'js>>(
ctx: Ctx<'js>,
mut objects: WritableStreamObjects<'js, W>,
) -> Result<WritableStreamObjects<'js, W>> {
objects
.stream
.writable_stream_mark_close_request_in_flight();
objects.controller.container.dequeue_value();
let (sink_close_promise, objects_class) = Self::close_algorithm(&ctx, objects)?;
objects = WritableStreamObjects::from_class(objects_class.clone());
objects
.controller
.writable_stream_default_controller_clear_algorithms();
upon_promise::<Value<'js>, ()>(ctx, sink_close_promise, |ctx, result| {
let objects = WritableStreamObjects::from_class(objects_class);
match result {
Ok(_) => {
WritableStream::writable_stream_finish_in_flight_close(objects)?;
},
Err(reason) => {
WritableStream::writable_stream_finish_in_flight_close_with_error(
ctx, objects, reason,
)?;
},
}
Ok(())
})?;
Ok(objects)
}
fn writable_stream_default_controller_process_write<W: WritableStreamWriter<'js>>(
ctx: Ctx<'js>,
mut objects: WritableStreamObjects<'js, W>,
chunk: Value<'js>,
) -> Result<WritableStreamObjects<'js, W>> {
objects
.stream
.writable_stream_mark_first_write_request_in_flight();
let (sink_write_promise, objects_class) = Self::write_algorithm(&ctx, objects, chunk)?;
upon_promise::<Value<'js>, ()>(ctx, sink_write_promise, {
let objects_class = objects_class.clone();
|ctx, result| {
let mut objects = WritableStreamObjects::from_class(objects_class).refresh_writer();
match result {
Ok(_) => {
objects.stream.writable_stream_finish_in_flight_write()?;
let state = &objects.stream.state;
objects.controller.container.dequeue_value();
if !objects.stream.writable_stream_close_queued_or_in_flight()
&& matches!(state, WritableStreamState::Writable)
{
let backpressure = objects
.controller
.writable_stream_default_controller_get_backpressure();
objects = WritableStream::writable_stream_update_backpressure(
ctx.clone(),
objects,
backpressure,
)?;
}
WritableStreamDefaultController::writable_stream_default_controller_advance_queue_if_needed(ctx, objects)?;
},
Err(reason) => {
if let WritableStreamState::Writable = objects.stream.state {
objects
.controller
.writable_stream_default_controller_clear_algorithms();
}
WritableStream::writable_stream_finish_in_flight_write_with_error(
ctx, objects, reason,
)?;
},
}
Ok(())
}
})?;
Ok(WritableStreamObjects::from_class(objects_class))
}
pub(super) fn error_steps(&mut self) {
self.reset_queue()
}
fn reset_queue(&mut self) {
self.container.queue.clear();
self.container.queue_total_size = 0.0;
}
pub(super) fn abort_steps<W: WritableStreamWriter<'js>>(
ctx: &Ctx<'js>,
mut objects: WritableStreamObjects<'js, W>,
reason: Value<'js>,
) -> Result<(Promise<'js>, WritableStreamObjects<'js, W>)> {
let (result, objects_class) = Self::abort_algorithm(ctx, objects, reason)?;
objects = WritableStreamObjects::from_class(objects_class);
objects
.controller
.writable_stream_default_controller_clear_algorithms();
Ok((result, objects))
}
fn strategy_size_algorithm(
ctx: Ctx<'js>,
objects: WritableStreamObjects<'js, WritableStreamDefaultWriterOwned<'js>>,
chunk: Value<'js>,
) -> (
Result<SizeValue<'js>>,
WritableStreamClassObjects<'js, WritableStreamDefaultWriterOwned<'js>>,
) {
let strategy_size_algorithm = objects
.controller
.strategy_size_algorithm
.clone()
.unwrap_or(SizeAlgorithm::AlwaysOne);
let objects_class = objects.into_inner();
(strategy_size_algorithm.call(ctx, chunk), objects_class)
}
fn start_algorithm(
ctx: Ctx<'js>,
objects: WritableStreamObjects<'js, UndefinedWriter>,
start_algorithm: WritableStartAlgorithm<'js>,
) -> Result<(Value<'js>, WritableStreamClassObjects<'js, UndefinedWriter>)> {
let objects_class = objects.into_inner();
Ok((
start_algorithm.call(ctx, objects_class.controller.clone())?,
objects_class,
))
}
fn write_algorithm<W: WritableStreamWriter<'js>>(
ctx: &Ctx<'js>,
objects: WritableStreamObjects<'js, W>,
chunk: Value<'js>,
) -> Result<(Promise<'js>, WritableStreamClassObjects<'js, W>)> {
let write_algorithm =
objects.controller.write_algorithm.clone().expect(
"write algorithm used after WritableStreamDefaultControllerClearAlgorithms",
);
let promise_primordials = objects.stream.promise_primordials.clone();
let objects_class = objects.into_inner();
Ok((
write_algorithm.call(
ctx,
&promise_primordials,
objects_class.controller.clone().clone(),
chunk,
)?,
objects_class,
))
}
fn close_algorithm<W: WritableStreamWriter<'js>>(
ctx: &Ctx<'js>,
objects: WritableStreamObjects<'js, W>,
) -> Result<(Promise<'js>, WritableStreamClassObjects<'js, W>)> {
let close_algorithm =
objects.controller.close_algorithm.clone().expect(
"close algorithm used after WritableStreamDefaultControllerClearAlgorithms",
);
let promise_primordials = objects.stream.promise_primordials.clone();
let objects_class = objects.into_inner();
Ok((
close_algorithm.call(ctx, &promise_primordials)?,
objects_class,
))
}
fn abort_algorithm<W: WritableStreamWriter<'js>>(
ctx: &Ctx<'js>,
objects: WritableStreamObjects<'js, W>,
reason: Value<'js>,
) -> Result<(Promise<'js>, WritableStreamClassObjects<'js, W>)> {
let abort_algorithm =
objects.controller.abort_algorithm.clone().expect(
"abort algorithm used after WritableStreamDefaultControllerClearAlgorithms",
);
let promise_primordials = objects.stream.promise_primordials.clone();
let objects_class = objects.into_inner();
Ok((
abort_algorithm.call(ctx, &promise_primordials, reason)?,
objects_class,
))
}
}
#[methods(rename_all = "camelCase")]
impl<'js> WritableStreamDefaultController<'js> {
#[qjs(get)]
pub fn constructor(ctx: Ctx<'js>) -> Result<Option<Constructor<'js>>> {
<WritableStreamDefaultController as JsClass>::constructor(&ctx)
}
#[qjs(constructor)]
fn new(ctx: Ctx<'js>) -> Result<Class<'js, Self>> {
Err(Exception::throw_type(&ctx, "Illegal constructor"))
}
#[qjs(get)]
fn signal(&self) -> Class<'js, AbortSignal<'js>> {
self.abort_controller.borrow().signal()
}
fn error(
ctx: Ctx<'js>,
controller: This<OwnedBorrowMut<'js, Self>>,
e: Opt<Value<'js>>,
) -> Result<()> {
let objects = WritableStreamObjects::from_controller(controller.0);
if !matches!(objects.stream.state, WritableStreamState::Writable) {
return Ok(());
}
Self::writable_stream_default_controller_error(
ctx.clone(),
objects.refresh_writer(),
e.0.unwrap_or_undefined(&ctx),
)?;
Ok(())
}
}
#[derive(Clone)]
pub(crate) enum WritableStartAlgorithm<'js> {
ReturnUndefined,
Function {
f: Function<'js>,
underlying_sink: Null<Undefined<Object<'js>>>,
},
Transform(Promise<'js>),
}
impl<'js> WritableStartAlgorithm<'js> {
fn call(
&self,
ctx: Ctx<'js>,
controller: WritableStreamDefaultControllerClass<'js>,
) -> Result<Value<'js>> {
match self {
WritableStartAlgorithm::ReturnUndefined => Ok(Value::new_undefined(ctx.clone())),
WritableStartAlgorithm::Function { f, underlying_sink } => {
f.call::<_, Value>((This(underlying_sink.clone()), controller))
},
WritableStartAlgorithm::Transform(promise) => Ok(promise.clone().into_value()),
}
}
}
#[derive(JsLifetime, Trace, Clone)]
pub(crate) enum WritableWriteAlgorithm<'js> {
ReturnPromiseUndefined,
Function {
f: Function<'js>,
underlying_sink: Null<Undefined<Object<'js>>>,
},
Transform {
stream: TransformStreamClass<'js>,
controller: TransformStreamDefaultControllerClass<'js>,
},
}
impl<'js> WritableWriteAlgorithm<'js> {
fn call(
&self,
ctx: &Ctx<'js>,
promise_primordials: &PromisePrimordials<'js>,
controller: WritableStreamDefaultControllerClass<'js>,
chunk: Value<'js>,
) -> Result<Promise<'js>> {
match self {
WritableWriteAlgorithm::ReturnPromiseUndefined => {
Ok(promise_primordials.promise_resolved_with_undefined.clone())
},
WritableWriteAlgorithm::Function { f, underlying_sink } => promise_resolved_with(
ctx,
promise_primordials,
f.call::<_, Value>((This(underlying_sink.clone()), chunk, controller)),
),
WritableWriteAlgorithm::Transform {
stream,
controller: ts_controller,
} => crate::stream_web::transform::stream::sink_write_algorithm(
ctx.clone(),
stream,
ts_controller,
chunk,
),
}
}
}
#[derive(JsLifetime, Trace, Clone)]
pub(crate) enum WritableCloseAlgorithm<'js> {
ReturnPromiseUndefined,
Function {
f: Function<'js>,
underlying_sink: Null<Undefined<Object<'js>>>,
},
Transform {
stream: TransformStreamClass<'js>,
controller: TransformStreamDefaultControllerClass<'js>,
},
}
impl<'js> WritableCloseAlgorithm<'js> {
fn call(
&self,
ctx: &Ctx<'js>,
promise_primordials: &PromisePrimordials<'js>,
) -> Result<Promise<'js>> {
match self {
WritableCloseAlgorithm::ReturnPromiseUndefined => {
Ok(promise_primordials.promise_resolved_with_undefined.clone())
},
WritableCloseAlgorithm::Function { f, underlying_sink } => promise_resolved_with(
ctx,
promise_primordials,
f.call::<_, Value>((This(underlying_sink.clone()),)),
),
WritableCloseAlgorithm::Transform { stream, controller } => {
crate::stream_web::transform::stream::sink_close_algorithm(ctx.clone(), stream, controller)
},
}
}
}
#[derive(JsLifetime, Trace, Clone)]
pub(crate) enum WritableAbortAlgorithm<'js> {
ReturnPromiseUndefined,
Function {
f: Function<'js>,
underlying_sink: Null<Undefined<Object<'js>>>,
},
Transform {
controller: TransformStreamDefaultControllerClass<'js>,
},
}
impl<'js> WritableAbortAlgorithm<'js> {
fn call(
&self,
ctx: &Ctx<'js>,
promise_primordials: &PromisePrimordials<'js>,
reason: Value<'js>,
) -> Result<Promise<'js>> {
match self {
WritableAbortAlgorithm::ReturnPromiseUndefined => {
Ok(promise_primordials.promise_resolved_with_undefined.clone())
},
WritableAbortAlgorithm::Function { f, underlying_sink } => promise_resolved_with(
ctx,
promise_primordials,
f.call::<_, Value>((This(underlying_sink.clone()), reason)),
),
WritableAbortAlgorithm::Transform { controller } => {
crate::stream_web::transform::stream::sink_abort_algorithm(ctx.clone(), controller, reason)
},
}
}
}
#[derive(Trace, Clone, JsLifetime)]
pub(crate) struct WritableStreamDefaultControllerPrimordials<'js> {
close_sentinel: Symbol<'js>,
}
impl<'js> Primordial<'js> for WritableStreamDefaultControllerPrimordials<'js> {
fn new(ctx: &Ctx<'js>) -> Result<Self>
where
Self: Sized,
{
Ok(Self {
close_sentinel: Symbol::new_global(ctx.clone(), "close sentinel")?,
})
}
}