use std::cell::RefCell;
use std::rc::Rc;
use js_sys::{SharedArrayBuffer, Uint8Array};
use wasm_bindgen::prelude::*;
use wasm_bindgen_futures::future_to_promise;
use crate::backend::SharedArrayBufferStorage;
use crate::header::HEADER_SIZE;
use crate::{Error, Options, Reader, Writer};
const MIN_POLL_MS: u32 = 1;
const MAX_POLL_MS: u32 = 4;
fn to_js_err(e: Error) -> JsValue {
JsValue::from_str(&e.to_string())
}
#[wasm_bindgen(getter_with_clone)]
pub struct ReadResult {
pub n: u32,
pub eof: bool,
}
#[wasm_bindgen(getter_with_clone)]
pub struct AsyncReadResult {
pub data: Uint8Array,
pub n: u32,
pub eof: bool,
}
#[wasm_bindgen(getter_with_clone)]
pub struct CreateWriterResult {
pub writer: WasmWriter,
pub sab: SharedArrayBuffer,
}
#[wasm_bindgen(js_name = createWriter)]
pub fn create_writer(capacity: u32) -> Result<CreateWriterResult, JsValue> {
let storage =
SharedArrayBufferStorage::new(HEADER_SIZE + capacity as u64).map_err(to_js_err)?;
let sab = storage.buffer();
let writer =
Writer::new(storage, capacity as u64, Options::default()).map_err(to_js_err)?;
Ok(CreateWriterResult {
writer: WasmWriter(Rc::new(RefCell::new(writer))),
sab,
})
}
#[wasm_bindgen(js_name = openReader)]
pub fn open_reader(sab: SharedArrayBuffer, capacity: u32) -> Result<WasmReader, JsValue> {
let storage = SharedArrayBufferStorage::wrap(sab).map_err(to_js_err)?;
let reader =
Reader::new(storage, capacity as u64, Options::default()).map_err(to_js_err)?;
Ok(WasmReader(Rc::new(RefCell::new(Some(reader)))))
}
#[derive(Clone)]
#[wasm_bindgen]
pub struct WasmWriter(Rc<RefCell<Writer<SharedArrayBufferStorage>>>);
#[wasm_bindgen]
impl WasmWriter {
pub fn try_write(&self, data: &[u8]) -> Result<u32, JsValue> {
self.0
.borrow_mut()
.try_write(data)
.map(|n| n as u32)
.map_err(to_js_err)
}
pub fn write(&self, data: Vec<u8>) -> js_sys::Promise {
let inner = self.0.clone();
future_to_promise(async move {
let mut written = 0usize;
let mut wait_ms = MIN_POLL_MS;
while written < data.len() {
let n = inner
.borrow_mut()
.try_write(&data[written..])
.map_err(to_js_err)?;
written += n;
if n > 0 {
wait_ms = MIN_POLL_MS;
continue;
}
gloo_timers::future::TimeoutFuture::new(wait_ms).await;
wait_ms = (wait_ms * 2).min(MAX_POLL_MS);
}
Ok(JsValue::from(written as u32))
})
}
pub fn close(&self) -> Result<(), JsValue> {
self.0.borrow_mut().close().map_err(to_js_err)
}
}
#[wasm_bindgen]
pub struct WasmReader(Rc<RefCell<Option<Reader<SharedArrayBufferStorage>>>>);
#[wasm_bindgen]
impl WasmReader {
pub fn try_read(&self, into: &mut [u8]) -> Result<ReadResult, JsValue> {
let mut guard = self.0.borrow_mut();
let reader = guard.as_mut().ok_or_else(|| to_js_err(Error::Closed))?;
match reader.try_read(into) {
Ok(n) => Ok(ReadResult {
n: n as u32,
eof: false,
}),
Err(Error::Eof) => Ok(ReadResult { n: 0, eof: true }),
Err(e) => Err(to_js_err(e)),
}
}
pub fn read(&self, len: u32) -> js_sys::Promise {
let inner = self.0.clone();
future_to_promise(async move {
let mut scratch = vec![0u8; len as usize];
let mut wait_ms = MIN_POLL_MS;
loop {
let outcome = {
let mut guard = inner.borrow_mut();
let reader = guard.as_mut().ok_or(Error::Closed)?;
reader.try_read(&mut scratch)
};
match outcome {
Ok(0) => {}
Ok(n) => {
let data = Uint8Array::new_with_length(n as u32);
data.copy_from(&scratch[..n]);
return Ok(JsValue::from(AsyncReadResult {
data,
n: n as u32,
eof: false,
}));
}
Err(Error::Eof) => {
return Ok(JsValue::from(AsyncReadResult {
data: Uint8Array::new_with_length(0),
n: 0,
eof: true,
}));
}
Err(e) => return Err(to_js_err(e)),
}
gloo_timers::future::TimeoutFuture::new(wait_ms).await;
wait_ms = (wait_ms * 2).min(MAX_POLL_MS);
}
})
}
pub fn close(&self) -> Result<(), JsValue> {
if let Some(reader) = self.0.borrow_mut().take() {
reader.close().map_err(to_js_err)?;
}
Ok(())
}
}
impl From<Error> for JsValue {
fn from(e: Error) -> JsValue {
to_js_err(e)
}
}