use std::{
cell::RefCell,
rc::Rc,
sync::atomic::{AtomicBool, Ordering},
};
use crate::abort::AbortSignal;
use crate::utils::{option::Undefined, result::ResultExt};
use rquickjs::{
class::{OwnedBorrow, Trace},
prelude::{OnceFn, This},
Class, Coerced, Ctx, Error, FromJs, Function, Promise, Result, Value,
};
use crate::stream_web::{
readable::{
controller::ReadableStreamControllerOwned,
default_reader::{
ReadableStreamDefaultReader, ReadableStreamDefaultReaderOwned,
ReadableStreamReadRequest,
},
objects::{
ReadableStreamClassObjects, ReadableStreamDefaultReaderObjects, ReadableStreamObjects,
},
reader::ReadableStreamReaderClass,
stream::{ReadableStream, ReadableStreamOwned, ReadableStreamState},
},
utils::{
promise::{
promise_resolved_with, upon_promise, upon_promise_fulfilment, PromisePrimordials,
ResolveablePromise,
},
UnwrapOrUndefined, ValueOrUndefined,
},
writable::{
WritableStream, WritableStreamClassObjects, WritableStreamDefaultWriter,
WritableStreamDefaultWriterOwned, WritableStreamObjects, WritableStreamOwned,
WritableStreamState,
},
};
impl<'js> ReadableStream<'js> {
pub(super) fn readable_stream_pipe_to(
ctx: Ctx<'js>,
source: ReadableStreamOwned<'js>,
dest: WritableStreamOwned<'js>,
prevent_close: bool,
prevent_abort: bool,
prevent_cancel: bool,
signal: Option<Class<'js, AbortSignal<'js>>>,
) -> Result<Promise<'js>> {
let (source_stored_error, source_closed) = match source.state {
ReadableStreamState::Errored(ref stored_error) => (Some(stored_error.clone()), false),
ReadableStreamState::Closed => (None, true),
_ => (None, false),
};
let dest_stored_error = dest.stored_error();
let dest_closing = dest.writable_stream_close_queued_or_in_flight()
|| matches!(dest.state, WritableStreamState::Closed);
let source_controller = source.controller.clone();
let dest_controller = dest
.controller
.clone()
.expect("pipeTo called on writable stream without controller");
let (mut source, reader) =
ReadableStreamReaderClass::acquire_readable_stream_default_reader(ctx.clone(), source)?;
let source_closed_promise = reader.borrow().generic.closed_promise.promise.clone();
let (dest, writer) =
WritableStreamDefaultWriter::acquire_writable_stream_default_writer(&ctx, dest)?;
let dest_closed_promise = writer.borrow().closed_promise.promise.clone();
source.disturbed = true;
let current_write = Rc::new(RefCell::new(
source
.promise_primordials
.promise_resolved_with_undefined
.clone(),
));
let promise_primordials = source.promise_primordials.clone();
let constructor_type_error = source.constructor_type_error.clone();
let mut pipe_to = PipeTo {
source_objects: ReadableStreamClassObjects {
stream: source.into_inner(),
controller: source_controller,
reader,
},
dest_objects: WritableStreamClassObjects {
stream: dest.into_inner(),
controller: dest_controller,
writer,
},
current_write,
shutting_down: Rc::new(AtomicBool::new(false)),
signal,
abort_callback: None,
promise: ResolveablePromise::new(&ctx)?,
promise_primordials: promise_primordials.clone(),
};
if let Some(signal) = &pipe_to.signal {
let abort_algorithm = {
let signal = signal.clone();
let pipe_to = pipe_to.clone();
move |ctx: Ctx<'js>| -> Result<()> {
let error = signal.borrow().reason().unwrap_or_undefined(&ctx);
let mut actions =
Vec::<Box<dyn FnOnce(Ctx<'js>) -> Result<Promise<'js>>>>::new();
if !prevent_abort {
let dest_objects = pipe_to.dest_objects.clone();
let error = error.clone();
actions.push(Box::new(move |ctx| {
let dest_objects = WritableStreamObjects::from_class(dest_objects);
if matches!(dest_objects.stream.state, WritableStreamState::Writable) {
let (promise, _) = WritableStream::writable_stream_abort(
ctx,
dest_objects,
Some(error.clone()),
)?;
Ok(promise)
} else {
Ok(dest_objects
.stream
.promise_primordials
.promise_resolved_with_undefined
.clone())
}
}));
}
if !prevent_cancel {
let source_objects = pipe_to.source_objects.clone();
let error = error.clone();
actions.push(Box::new(move |ctx| {
let source_objects = ReadableStreamObjects::from_class(source_objects);
if let ReadableStreamState::Readable = source_objects.stream.state {
let (promise, _) = ReadableStream::readable_stream_cancel(
ctx,
source_objects,
error.clone(),
)?;
Ok(promise)
} else {
Ok(source_objects
.stream
.promise_primordials
.promise_resolved_with_undefined
.clone())
}
}));
}
pipe_to.shutdown_with_action(
ctx,
move |ctx| {
let promises: Vec<Promise<'js>> = actions
.into_iter()
.map(|action| action(ctx.clone()))
.collect::<Result<Vec<_>>>()?;
let all_promises: Promise<'js> =
promise_primordials.promise_all.call((
This(promise_primordials.promise_constructor.clone()),
promises,
))?;
Ok(all_promises)
},
Some(error),
)
}
};
{
let signal = signal.borrow();
if signal.aborted {
abort_algorithm(ctx.clone())?;
return Ok(pipe_to.promise.promise);
}
}
let abort_callback = pipe_to
.abort_callback
.insert(Function::new(ctx.clone(), OnceFn::new(abort_algorithm))?);
AbortSignal::set_on_abort(This(signal.clone()), ctx.clone(), abort_callback.clone())?;
}
PipeTo::is_or_becomes_errored(
ctx.clone(),
source_stored_error,
source_closed_promise.clone(),
{
let pipe_to = pipe_to.clone();
move |ctx, stored_error| {
if !prevent_abort {
pipe_to.shutdown_with_action(
ctx,
{
let pipe_to = pipe_to.clone();
let stored_error = stored_error.clone();
move |ctx| {
let dest_objects = WritableStreamObjects::from_class(
pipe_to.dest_objects.clone(),
);
let (promise, _) = WritableStream::writable_stream_abort(
ctx,
dest_objects,
Some(stored_error),
)?;
Ok(promise)
}
},
Some(stored_error),
)
} else {
pipe_to.shutdown(ctx, Some(stored_error))
}
}
},
)?;
PipeTo::is_or_becomes_errored(ctx.clone(), dest_stored_error, dest_closed_promise, {
let pipe_to = pipe_to.clone();
move |ctx, stored_error| {
if !prevent_cancel {
pipe_to.shutdown_with_action(
ctx,
{
let pipe_to = pipe_to.clone();
let stored_error = stored_error.clone();
move |ctx| {
let source_objects = ReadableStreamObjects::from_class(
pipe_to.source_objects.clone(),
);
let (promise, _) = ReadableStream::readable_stream_cancel(
ctx,
source_objects,
stored_error,
)?;
Ok(promise)
}
},
Some(stored_error),
)
} else {
pipe_to.shutdown(ctx, Some(stored_error))
}
}
})?;
PipeTo::is_or_becomes_closed(ctx.clone(), source_closed, source_closed_promise, {
let pipe_to = pipe_to.clone();
move |ctx| {
if !prevent_close {
pipe_to.shutdown_with_action(
ctx,
{
let pipe_to = pipe_to.clone();
move |ctx| {
let dest_objects = WritableStreamObjects::from_class(pipe_to.dest_objects);
WritableStreamDefaultWriter::writable_stream_default_writer_close_with_error_propagation(ctx, dest_objects)
}
},
None,
)
} else {
pipe_to.shutdown(ctx, None)
}
}
})?;
if dest_closing {
let dest_closed: Value<'js> = constructor_type_error.call((
"the destination writable stream closed before all data could be piped to it",
))?;
if !prevent_cancel {
pipe_to.shutdown_with_action(
ctx.clone(),
{
let pipe_to = pipe_to.clone();
let dest_closed = dest_closed.clone();
move |ctx| {
let source_objects =
ReadableStreamObjects::from_class(pipe_to.source_objects.clone());
let (promise, _) = ReadableStream::readable_stream_cancel(
ctx,
source_objects,
dest_closed,
)?;
Ok(promise)
}
},
Some(dest_closed),
)?;
} else {
pipe_to.shutdown(ctx.clone(), Some(dest_closed))?;
}
}
let result_promise = pipe_to.promise.promise.clone();
let pipe_loop_promise = pipe_to.pipe_loop(ctx)?;
pipe_loop_promise.set_is_handled()?;
Ok(result_promise)
}
}
#[derive(Clone)]
struct PipeTo<'js> {
source_objects: ReadableStreamClassObjects<
'js,
ReadableStreamControllerOwned<'js>,
ReadableStreamDefaultReaderOwned<'js>,
>,
dest_objects: WritableStreamClassObjects<'js, WritableStreamDefaultWriterOwned<'js>>,
current_write: Rc<RefCell<Promise<'js>>>,
shutting_down: Rc<AtomicBool>,
signal: Option<Class<'js, AbortSignal<'js>>>,
abort_callback: Option<Function<'js>>,
promise: ResolveablePromise<'js>,
promise_primordials: PromisePrimordials<'js>,
}
impl<'js> PipeTo<'js> {
fn pipe_loop(self, ctx: Ctx<'js>) -> Result<ResolveablePromise<'js>> {
let loop_promise = ResolveablePromise::new(&ctx)?;
self.next(ctx, false, loop_promise.clone())?;
Ok(loop_promise)
}
fn next(&self, ctx: Ctx<'js>, done: bool, loop_promise: ResolveablePromise<'js>) -> Result<()> {
if done {
loop_promise.resolve_undefined()?
} else {
let pipe_step_promise = self.pipe_step(ctx.clone())?;
upon_promise(ctx, pipe_step_promise, {
{
let pipe_to = self.clone();
move |ctx, result| match result {
Ok(done) => pipe_to.next(ctx, done, loop_promise),
Err(err) => loop_promise.reject(err),
}
}
})?;
}
Ok(())
}
fn pipe_step(&self, ctx: Ctx<'js>) -> Result<Promise<'js>> {
if self.shutting_down.load(Ordering::Acquire) {
return promise_resolved_with(
&ctx,
&self.promise_primordials,
Ok(Value::new_bool(ctx.clone(), true)),
);
}
let writer_ready = self
.dest_objects
.writer
.borrow()
.ready_promise
.promise
.clone();
upon_promise_fulfilment(ctx, writer_ready, {
let current_write = self.current_write.clone();
let source_objects = self.source_objects.clone();
let dest_objects = self.dest_objects.clone();
move |ctx: Ctx<'js>, ()| -> Result<Promise<'js>> {
let read_promise = ResolveablePromise::new(&ctx)?;
struct ReadRequest<'js> {
dest_objects:
WritableStreamClassObjects<'js, WritableStreamDefaultWriterOwned<'js>>,
current_write: Rc<RefCell<Promise<'js>>>,
read_promise: ResolveablePromise<'js>,
}
impl<'js> Trace<'js> for ReadRequest<'js> {
fn trace<'a>(&self, tracer: rquickjs::class::Tracer<'a, 'js>) {
self.current_write.as_ref().borrow().trace(tracer);
self.read_promise.trace(tracer);
}
}
impl<'js> ReadableStreamReadRequest<'js> for ReadRequest<'js> {
fn chunk_steps(
&self,
objects: ReadableStreamDefaultReaderObjects<'js>,
chunk: Value<'js>,
) -> Result<ReadableStreamDefaultReaderObjects<'js>> {
let ctx = chunk.ctx().clone();
let objects = objects.into_inner();
let dest_objects =
WritableStreamObjects::from_class(self.dest_objects.clone());
let write_promise =
WritableStreamDefaultWriter::writable_stream_default_writer_write(
ctx.clone(),
dest_objects,
chunk,
)?;
let write_promise: Promise<'js> = write_promise.catch()?.call((
This(write_promise.clone()),
Function::new(ctx.clone(), || {}),
))?;
self.current_write.replace(write_promise);
self.read_promise
.resolve(Value::new_bool(ctx.clone(), false))?;
Ok(ReadableStreamObjects::from_class(objects))
}
fn close_steps(
&self,
ctx: &Ctx<'js>,
objects: ReadableStreamDefaultReaderObjects<'js>,
) -> Result<ReadableStreamDefaultReaderObjects<'js>> {
self.read_promise
.resolve(Value::new_bool(ctx.clone(), true))?;
Ok(objects)
}
fn error_steps(
&self,
objects: ReadableStreamDefaultReaderObjects<'js>,
reason: Value<'js>,
) -> Result<ReadableStreamDefaultReaderObjects<'js>> {
self.read_promise.reject(reason)?;
Ok(objects)
}
}
let objects = ReadableStreamObjects::from_class(source_objects);
let promise = read_promise.promise.clone();
ReadableStreamDefaultReader::readable_stream_default_reader_read(
&ctx,
objects,
ReadRequest {
current_write,
read_promise,
dest_objects,
},
)?;
Ok(promise)
}
})
}
fn is_or_becomes_errored(
ctx: Ctx<'js>,
stored_error: Option<Value<'js>>,
promise: Promise<'js>,
action: impl FnOnce(Ctx<'js>, Value<'js>) -> Result<()> + 'js,
) -> Result<()> {
if let Some(stored_error) = stored_error {
action(ctx, stored_error)
} else {
promise.catch()?.call((
This(promise.clone()),
Function::new(ctx.clone(), OnceFn::new(action)),
))
}
}
fn is_or_becomes_closed(
ctx: Ctx<'js>,
already_closed: bool,
promise: Promise<'js>,
action: impl FnOnce(Ctx<'js>) -> Result<()> + 'js,
) -> Result<()> {
if already_closed {
action(ctx)?;
} else {
upon_promise_fulfilment(ctx, promise, |ctx, ()| action(ctx))?;
}
Ok(())
}
fn shutdown_with_action(
&self,
ctx: Ctx<'js>,
action: impl FnOnce(Ctx<'js>) -> Result<Promise<'js>> + 'js,
original_error: Option<Value<'js>>,
) -> Result<()> {
if self.shutting_down.swap(true, Ordering::AcqRel) {
return Ok(());
}
let do_the_rest = {
let pipe_to = self.clone();
move |ctx: Ctx<'js>| -> Result<()> {
let action_promise = action(ctx.clone())?;
upon_promise(ctx, action_promise, move |ctx, result| match result {
Ok(()) => pipe_to.finalize(ctx, original_error),
Err(new_error) => pipe_to.finalize(ctx, Some(new_error)),
})?;
Ok(())
}
};
let writable = {
let dest_stream = OwnedBorrow::from_class(self.dest_objects.stream.clone());
matches!(dest_stream.state, WritableStreamState::Writable)
&& !dest_stream.writable_stream_close_queued_or_in_flight()
};
if writable {
let wait_promise =
Self::wait_for_writes_to_finish(ctx.clone(), self.current_write.clone())?;
upon_promise_fulfilment(ctx, wait_promise, |ctx: Ctx<'js>, ()| do_the_rest(ctx))?;
} else {
do_the_rest(ctx)?
}
Ok(())
}
fn shutdown(&self, ctx: Ctx<'js>, error: Option<Value<'js>>) -> Result<()> {
if self.shutting_down.swap(true, Ordering::AcqRel) {
return Ok(());
}
let writable = {
let dest_stream = OwnedBorrow::from_class(self.dest_objects.stream.clone());
matches!(dest_stream.state, WritableStreamState::Writable)
&& !dest_stream.writable_stream_close_queued_or_in_flight()
};
if writable {
let wait_promise =
Self::wait_for_writes_to_finish(ctx.clone(), self.current_write.clone())?;
let pipe_to = self.clone();
upon_promise_fulfilment(ctx, wait_promise, move |ctx, ()| {
pipe_to.finalize(ctx, error)
})?;
} else {
self.finalize(ctx, error)?;
}
Ok(())
}
fn wait_for_writes_to_finish(
ctx: Ctx<'js>,
current_write: Rc<RefCell<Promise<'js>>>,
) -> Result<Promise<'js>> {
let old_current_write: Promise<'js> = current_write.as_ref().borrow().clone();
upon_promise_fulfilment(
ctx,
old_current_write.clone(),
move |ctx: Ctx<'js>, ()| -> Result<Undefined<Promise<'js>>> {
if !old_current_write.eq(¤t_write.as_ref().borrow()) {
Ok(Undefined(Some(Self::wait_for_writes_to_finish(
ctx,
current_write,
)?)))
} else {
Ok(Undefined(None))
}
},
)
}
fn finalize(&self, ctx: Ctx<'js>, error: Option<Value<'js>>) -> Result<()> {
let source_objects = ReadableStreamObjects::from_class(self.source_objects.clone());
let dest_objects = WritableStreamObjects::from_class(self.dest_objects.clone());
WritableStreamDefaultWriter::writable_stream_default_writer_release(dest_objects)?;
ReadableStreamDefaultReader::readable_stream_default_reader_release(source_objects)?;
if let (Some(signal), Some(abort_callback)) = (&self.signal, &self.abort_callback) {
AbortSignal::remove_on_abort(
This(signal.clone()),
ctx.clone(),
abort_callback.clone(),
)?;
}
if let Some(error) = error {
self.promise.reject(error)
} else {
self.promise.resolve_undefined()
}
}
}
#[derive(Default)]
pub struct StreamPipeOptions<'js> {
pub prevent_close: bool,
pub prevent_abort: bool,
pub prevent_cancel: bool,
pub signal: Option<Class<'js, AbortSignal<'js>>>,
}
impl<'js> FromJs<'js> for StreamPipeOptions<'js> {
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 get_bool = |key| {
Result::Ok(
obj.get_value_or_undefined::<_, Coerced<bool>>(key)?
.map(|b| b.0)
.unwrap_or(false),
) };
let prevent_abort = get_bool("preventAbort")?;
let prevent_cancel = get_bool("preventCancel")?;
let prevent_close = get_bool("preventClose")?;
let signal = match obj.get_value_or_undefined::<_, Value<'js>>("signal")? {
Some(signal) => Some(
Class::<AbortSignal>::from_js(ctx, signal)
.or_throw_type(ctx, "Invalid signal argument")?,
),
None => None,
};
Ok(Self {
prevent_close,
prevent_abort,
prevent_cancel,
signal,
})
}
}