use std::sync::Arc;
use ferrijs_fetch::ByteStream;
use ferrijs_std::context::CtxExtension;
use ferrijs_std::stream_web::utils::promise::{PromisePrimordials, ResolveablePromise};
use ferrijs_std::stream_web::{
CancelAlgorithm, PullAlgorithm, ReadableStream, ReadableStreamControllerClass, ReadableStreamDefaultControllerClass,
readable_stream_default_controller_close_stream, readable_stream_default_controller_enqueue_value,
readable_stream_default_controller_error_stream,
};
use ferrijs_std::utils::primordials::Primordial;
use futures::StreamExt as _;
use rquickjs::{CatchResultExt as _, Class, Ctx, TypedArray, Value};
use tokio::sync::Mutex as AsyncMutex;
pub type NetBody = Arc<AsyncMutex<Option<ByteStream>>>;
const NET_CHUNK_TIMEOUT: std::time::Duration = std::time::Duration::from_mins(2);
fn default_controller<'js>(
ctx: &Ctx<'js>,
controller: ReadableStreamControllerClass<'js>,
) -> rquickjs::Result<ReadableStreamDefaultControllerClass<'js>> {
match controller {
ReadableStreamControllerClass::ReadableStreamDefaultController(c) => Ok(c),
_ => Err(rquickjs::Exception::throw_type(
ctx,
"expected a default ReadableStream controller",
)),
}
}
fn resolved_undefined<'js>(ctx: &Ctx<'js>) -> rquickjs::Result<rquickjs::Promise<'js>> {
Ok(PromisePrimordials::get(ctx)?.promise_resolved_with_undefined.clone())
}
pub fn from_bytes<'js>(ctx: &Ctx<'js>, bytes: Vec<u8>) -> rquickjs::Result<Class<'js, ReadableStream<'js>>> {
let pull = PullAlgorithm::from_fn_once(move |ctx: Ctx<'js>, controller| {
let ctrl = default_controller(&ctx, controller)?;
let chunk = TypedArray::<u8>::new(ctx.clone(), bytes)?.into_value();
readable_stream_default_controller_enqueue_value(ctx.clone(), ctrl.clone(), chunk)?;
readable_stream_default_controller_close_stream(ctx.clone(), ctrl)?;
resolved_undefined(&ctx)
});
ReadableStream::from_pull_algorithm(ctx.clone(), pull, CancelAlgorithm::ReturnPromiseUndefined)
}
pub fn to_byte_stream<'js>(ctx: &Ctx<'js>, stream: Class<'js, ReadableStream<'js>>) -> rquickjs::Result<ByteStream> {
let (tx, rx) = tokio::sync::mpsc::channel::<Result<Vec<u8>, String>>(1);
let reader = {
let obj = stream
.clone()
.into_value()
.into_object()
.ok_or_else(|| rquickjs::Error::new_from_js_message("fetch", "TypeError", "body is not a ReadableStream"))?;
obj
.get::<_, rquickjs::Function<'js>>("getReader")?
.call::<_, rquickjs::Object<'js>>((rquickjs::function::This(obj),))?
};
let pump_ctx = ctx.clone();
ctx.spawn_exit_simple(async move {
let read: rquickjs::Function<'js> = reader.get("read")?;
loop {
let step: rquickjs::Promise<'js> = read.call((rquickjs::function::This(reader.clone()),))?;
let outcome = step.into_future::<rquickjs::Object<'js>>().await.catch(&pump_ctx);
let message = match outcome {
Ok(res) => {
if res.get::<_, bool>("done").unwrap_or(false) {
return Ok(());
}
Ok(chunk_bytes(&res.get::<_, Value<'js>>("value")?))
},
Err(e) => Err(crate::ScriptError::from_caught(&pump_ctx, e, "").message),
};
let failed = message.is_err();
if tx.send(message).await.is_err() || failed {
return Ok(());
}
}
});
Ok(ferrijs_fetch::channel_stream(rx))
}
fn chunk_bytes(v: &Value<'_>) -> Vec<u8> {
if let Some(s) = v.as_string().and_then(|s| s.to_string().ok()) {
return s.into_bytes();
}
if let Ok(ta) = TypedArray::<u8>::from_value(v.clone()) {
let bytes: &[u8] = ta.as_ref();
return bytes.to_vec();
}
if let Some(ab) = rquickjs::ArrayBuffer::from_value(v.clone())
&& let Some(bytes) = ab.as_bytes()
{
return bytes.to_vec();
}
Vec::new()
}
pub fn from_net<'js>(ctx: &Ctx<'js>, net: NetBody) -> rquickjs::Result<Class<'js, ReadableStream<'js>>> {
let pull_net = net.clone();
let pull = PullAlgorithm::from_fn(move |ctx: Ctx<'js>, controller| {
let ctrl = default_controller(&ctx, controller)?;
let net = pull_net.clone();
let resolveable = ResolveablePromise::new(&ctx)?;
let promise = resolveable.promise.clone();
let ctx2 = ctx.clone();
ctx.spawn_exit_simple(async move {
let mut guard = net.lock().await;
match guard.as_mut() {
None => readable_stream_default_controller_close_stream(ctx2, ctrl)?,
Some(body) => match tokio::time::timeout(NET_CHUNK_TIMEOUT, body.next()).await {
Ok(Some(Ok(bytes))) => {
let chunk = TypedArray::<u8>::new(ctx2.clone(), bytes.to_vec())?.into_value();
readable_stream_default_controller_enqueue_value(ctx2, ctrl, chunk)?;
},
Ok(None) => {
*guard = None;
readable_stream_default_controller_close_stream(ctx2, ctrl)?;
},
Ok(Some(Err(e))) => {
*guard = None;
let err = rquickjs::String::from_str(ctx2, &e.to_string())?.into_value();
readable_stream_default_controller_error_stream(ctrl, err)?;
},
Err(_) => {
*guard = None;
let err = rquickjs::String::from_str(ctx2, "body read timed out: no chunk within 120s")?.into_value();
readable_stream_default_controller_error_stream(ctrl, err)?;
},
},
}
resolveable.resolve_undefined()?;
Ok(())
});
Ok(promise)
});
let cancel = CancelAlgorithm::from_fn(move |reason: Value<'js>| {
if let Ok(mut g) = net.try_lock() {
*g = None;
}
resolved_undefined(reason.ctx())
});
ReadableStream::from_pull_algorithm(ctx.clone(), pull, cancel)
}