use crate::utils::option::Null;
use rquickjs::{
class::{JsClass, OwnedBorrow, OwnedBorrowMut, Trace},
function::Constructor,
prelude::{Opt, This},
Class, Ctx, Exception, JsLifetime, Promise, Result, Value,
};
use crate::stream_web::{
utils::{
promise::{
promise_rejected_with, promise_rejected_with_constructor, PromisePrimordials,
ResolveablePromise,
},
UnwrapOrUndefined,
},
writable::{
default_controller::WritableStreamDefaultController,
objects::WritableStreamObjects,
stream::{WritableStream, WritableStreamOwned, WritableStreamState},
writer::WritableStreamWriter,
},
};
#[rquickjs::class]
#[derive(JsLifetime)]
pub(crate) struct WritableStreamDefaultWriter<'js> {
pub(crate) ready_promise: ResolveablePromise<'js>,
pub(crate) closed_promise: ResolveablePromise<'js>,
pub(super) stream: Option<Class<'js, WritableStream<'js>>>,
constructor_type_error: Constructor<'js>,
promise_primordials: PromisePrimordials<'js>,
}
impl<'js> Trace<'js> for WritableStreamDefaultWriter<'js> {
fn trace<'a>(&self, tracer: rquickjs::class::Tracer<'a, 'js>) {
self.ready_promise.trace(tracer);
self.closed_promise.trace(tracer);
self.stream.trace(tracer);
self.constructor_type_error.trace(tracer);
self.promise_primordials.trace(tracer);
}
}
pub(crate) type WritableStreamDefaultWriterClass<'js> =
Class<'js, WritableStreamDefaultWriter<'js>>;
pub(crate) type WritableStreamDefaultWriterOwned<'js> =
OwnedBorrowMut<'js, WritableStreamDefaultWriter<'js>>;
#[rquickjs::methods(rename_all = "camelCase")]
impl<'js> WritableStreamDefaultWriter<'js> {
#[qjs(get)]
pub fn constructor(ctx: Ctx<'js>) -> Result<Option<Constructor<'js>>> {
<WritableStreamDefaultWriter as JsClass>::constructor(&ctx)
}
#[qjs(constructor)]
fn new(ctx: Ctx<'js>, stream: WritableStreamOwned<'js>) -> Result<Class<'js, Self>> {
let (_, writer) = Self::set_up_writable_stream_default_writer(&ctx, stream)?;
Ok(writer)
}
#[qjs(get)]
fn closed(writer: This<OwnedBorrowMut<'js, Self>>) -> Promise<'js> {
writer.0.closed_promise.promise.clone()
}
#[qjs(get)]
fn desired_size(ctx: Ctx<'js>, writer: This<OwnedBorrowMut<'js, Self>>) -> Result<Null<f64>> {
match writer.0.stream {
None => Err(Exception::throw_type(
&ctx,
"Cannot desiredSize a stream using a released writer",
)),
Some(ref stream) => {
Self::writable_stream_default_writer_get_desired_size(&OwnedBorrowMut::from_class(
stream.clone(),
))
},
}
}
#[qjs(get)]
fn ready(writer: This<OwnedBorrowMut<'js, Self>>) -> Promise<'js> {
writer.0.ready_promise.promise.clone()
}
fn abort(
ctx: Ctx<'js>,
writer: This<OwnedBorrowMut<'js, Self>>,
reason: Opt<Value<'js>>,
) -> Result<Promise<'js>> {
if writer.0.stream.is_none() {
promise_rejected_with_constructor(
&writer.constructor_type_error,
&writer.promise_primordials,
"Cannot abort a stream using a released writer",
)
} else {
let objects = WritableStreamObjects::from_writer(writer.0);
Self::writable_stream_default_writer_abort(ctx.clone(), objects, reason.0)
}
}
fn close(ctx: Ctx<'js>, writer: This<OwnedBorrowMut<'js, Self>>) -> Result<Promise<'js>> {
if writer.0.stream.is_none() {
promise_rejected_with_constructor(
&writer.constructor_type_error,
&writer.promise_primordials,
"Cannot close a stream using a released writer",
)
} else {
let objects = WritableStreamObjects::from_writer(writer.0);
if objects.stream.writable_stream_close_queued_or_in_flight() {
return promise_rejected_with_constructor(
&objects.writer.constructor_type_error,
&objects.writer.promise_primordials,
"Cannot close an already-closing",
);
}
Self::writable_stream_default_writer_close(ctx, objects)
}
}
fn release_lock(writer: This<OwnedBorrowMut<'js, Self>>) -> Result<()> {
if writer.0.stream.is_none() {
Ok(())
} else {
let objects = WritableStreamObjects::from_writer(writer.0);
Self::writable_stream_default_writer_release(objects)
}
}
fn write(
ctx: Ctx<'js>,
writer: This<OwnedBorrowMut<'js, Self>>,
chunk: Opt<Value<'js>>,
) -> Result<Promise<'js>> {
if writer.0.stream.is_none() {
promise_rejected_with_constructor(
&writer.constructor_type_error,
&writer.promise_primordials,
"Cannot write a stream using a released writer",
)
} else {
let objects = WritableStreamObjects::from_writer(writer.0);
Self::writable_stream_default_writer_write(
ctx.clone(),
objects,
chunk.0.unwrap_or_undefined(&ctx),
)
}
}
}
impl<'js> WritableStreamDefaultWriter<'js> {
pub(crate) fn acquire_writable_stream_default_writer(
ctx: &Ctx<'js>,
stream: WritableStreamOwned<'js>,
) -> Result<(WritableStreamOwned<'js>, Class<'js, Self>)> {
Self::set_up_writable_stream_default_writer(ctx, stream)
}
pub(super) fn set_up_writable_stream_default_writer(
ctx: &Ctx<'js>,
mut stream: WritableStreamOwned<'js>,
) -> Result<(WritableStreamOwned<'js>, Class<'js, Self>)> {
if stream.is_writable_stream_locked() {
return Err(Exception::throw_type(
ctx,
"This stream has already been locked for exclusive writing by another writer",
));
}
let promise_primordials = stream.promise_primordials.clone();
let constructor_type_error = stream.constructor_type_error.clone();
let stream_class = stream.into_inner();
stream = OwnedBorrowMut::from_class(stream_class.clone());
let (ready_promise, closed_promise) = match stream.state {
WritableStreamState::Writable => {
let ready_promise =
if !stream.writable_stream_close_queued_or_in_flight() && stream.backpressure {
ResolveablePromise::new(ctx)?
} else {
ResolveablePromise::resolved_with_undefined(&stream.promise_primordials)
};
(ready_promise, ResolveablePromise::new(ctx)?)
},
WritableStreamState::Erroring(ref stored_error) => {
let ready_promise = ResolveablePromise::rejected_with(
&stream.promise_primordials,
stored_error.clone(),
)?;
ready_promise.set_is_handled()?;
(ready_promise, ResolveablePromise::new(ctx)?)
},
WritableStreamState::Closed => {
let promise =
ResolveablePromise::resolved_with_undefined(&stream.promise_primordials);
(promise.clone(), promise)
},
WritableStreamState::Errored(ref stored_error) => {
let promise = ResolveablePromise::rejected_with(
&stream.promise_primordials,
stored_error.clone(),
)?;
promise.set_is_handled()?;
(promise.clone(), promise)
},
};
let writer = Self {
ready_promise,
closed_promise,
stream: Some(stream_class),
promise_primordials,
constructor_type_error,
};
let writer = Class::instance(ctx.clone(), writer)?;
stream.writer = Some(writer.clone());
Ok((stream, writer))
}
pub(super) fn writable_stream_default_writer_ensure_ready_promise_rejected(
&mut self,
promise_primordials: &PromisePrimordials<'js>,
error: Value<'js>,
) -> Result<()> {
if self.ready_promise.is_pending() {
self.ready_promise.reject(error)?;
} else {
self.ready_promise = ResolveablePromise::rejected_with(promise_primordials, error)?;
}
self.ready_promise.set_is_handled()?;
Ok(())
}
pub(super) fn writable_stream_default_writer_ensure_closed_promise_rejected(
&mut self,
promise_primordials: &PromisePrimordials<'js>,
error: Value<'js>,
) -> Result<()> {
if self.closed_promise.is_pending() {
self.closed_promise.reject(error)?;
} else {
self.closed_promise = ResolveablePromise::rejected_with(promise_primordials, error)?;
}
self.closed_promise.set_is_handled()?;
Ok(())
}
pub(super) fn writable_stream_default_writer_get_desired_size(
stream: &WritableStream<'js>,
) -> Result<Null<f64>> {
if matches!(
stream.state,
WritableStreamState::Errored(_) | WritableStreamState::Erroring(_)
) {
return Ok(Null(None));
}
if matches!(stream.state, WritableStreamState::Closed) {
return Ok(Null(Some(0.0)));
}
let controller = OwnedBorrow::from_class(
stream
.controller
.clone()
.expect("Stream in state writable must have a controller"),
);
Ok(Null(Some(
controller.writable_stream_default_controller_get_desired_size(),
)))
}
fn writable_stream_default_writer_abort(
ctx: Ctx<'js>,
objects: WritableStreamObjects<'js, OwnedBorrowMut<'js, Self>>,
reason: Option<Value<'js>>,
) -> Result<Promise<'js>> {
let (promise, _) = WritableStream::writable_stream_abort(ctx, objects, reason)?;
Ok(promise)
}
fn writable_stream_default_writer_close(
ctx: Ctx<'js>,
objects: WritableStreamObjects<'js, OwnedBorrowMut<'js, Self>>,
) -> Result<Promise<'js>> {
let (promise, _) = WritableStream::writable_stream_close(ctx, objects)?;
Ok(promise)
}
pub(crate) fn writable_stream_default_writer_close_with_error_propagation(
ctx: Ctx<'js>,
objects: WritableStreamObjects<'js, OwnedBorrowMut<'js, Self>>,
) -> Result<Promise<'js>> {
if objects.stream.writable_stream_close_queued_or_in_flight()
|| matches!(objects.stream.state, WritableStreamState::Closed)
{
return Ok(objects
.stream
.promise_primordials
.promise_resolved_with_undefined
.clone());
}
if let WritableStreamState::Errored(ref stored_error) = objects.stream.state {
return promise_rejected_with(
&objects.stream.promise_primordials,
stored_error.clone(),
);
}
Self::writable_stream_default_writer_close(ctx, objects)
}
pub(crate) fn writable_stream_default_writer_release(
mut objects: WritableStreamObjects<'js, OwnedBorrowMut<'js, Self>>,
) -> Result<()> {
let released_error: Value = objects.stream.constructor_type_error.call((
"Writer was released and can no longer be used to monitor the stream's closedness",
))?;
objects
.writer
.writable_stream_default_writer_ensure_ready_promise_rejected(
&objects.stream.promise_primordials,
released_error.clone(),
)?;
objects
.writer
.writable_stream_default_writer_ensure_closed_promise_rejected(
&objects.stream.promise_primordials,
released_error,
)?;
objects.stream.writer = None;
objects.writer.stream = None;
Ok(())
}
pub(crate) fn writable_stream_default_writer_write(
ctx: Ctx<'js>,
objects: WritableStreamObjects<'js, OwnedBorrowMut<'js, Self>>,
chunk: Value<'js>,
) -> Result<Promise<'js>> {
let (chunk_size, mut objects) =
WritableStreamDefaultController::writable_stream_default_controller_get_chunk_size(
ctx.clone(),
objects,
chunk.clone(),
)?;
let stream_class = objects.stream.into_inner();
objects.stream = OwnedBorrowMut::from_class(stream_class.clone());
if objects.writer.stream != Some(stream_class) {
return promise_rejected_with_constructor(
&objects.stream.constructor_type_error,
&objects.stream.promise_primordials,
"Cannot write to a stream using a released writer",
);
}
if let WritableStreamState::Errored(ref stored_error) = objects.stream.state {
return promise_rejected_with(
&objects.stream.promise_primordials,
stored_error.clone(),
);
}
if objects.stream.writable_stream_close_queued_or_in_flight()
|| matches!(objects.stream.state, WritableStreamState::Closed)
{
return promise_rejected_with_constructor(
&objects.stream.constructor_type_error,
&objects.stream.promise_primordials,
"The stream is closing or closed and cannot be written to",
);
}
if let WritableStreamState::Erroring(ref stored_error) = objects.stream.state {
return promise_rejected_with(
&objects.stream.promise_primordials,
stored_error.clone(),
);
}
let promise = objects.stream.writable_stream_add_write_request(&ctx);
WritableStreamDefaultController::writable_stream_default_controller_write(
ctx, objects, chunk, chunk_size,
)?;
promise
}
}
impl<'js> WritableStreamWriter<'js> for WritableStreamDefaultWriterOwned<'js> {
type Class = WritableStreamDefaultWriterClass<'js>;
fn with_writer<C>(
self,
ctx: C,
default: impl FnOnce(
C,
WritableStreamDefaultWriterOwned<'js>,
) -> Result<(C, WritableStreamDefaultWriterOwned<'js>)>,
_: impl FnOnce(C) -> Result<C>,
) -> Result<(C, Self)> {
default(ctx, self)
}
fn into_inner(self) -> Self::Class {
self.into_inner()
}
fn from_class(class: Self::Class) -> Self {
OwnedBorrowMut::from_class(class)
}
}