use std::cell::{Cell, RefCell};
use std::collections::HashMap;
use std::rc::Rc;
use async_trait::async_trait;
use js_sys::{Function, Object, Promise, Reflect, Uint8Array};
use wasm_bindgen::prelude::*;
use wasm_bindgen::JsCast;
use wasm_bindgen_futures::JsFuture;
use web_sys::{
File, FileSystemDirectoryHandle, FileSystemFileHandle, FileSystemGetDirectoryOptions,
FileSystemGetFileOptions, FileSystemHandle, FileSystemHandleKind, FileSystemRemoveOptions,
FileSystemWritableFileStream, MessageEvent, Worker,
};
use super::{DirEntry, EntryKind, Filesystem, Metadata, WalkEntry};
use crate::error::{Error, Result};
const WRITE_TIMEOUT_MS: u32 = 15_000;
const MAX_WALK_ENTRIES: usize = 200_000;
#[derive(Debug, Clone, Default)]
pub struct OpfsFilesystem {
root: Rc<RefCell<Option<FileSystemDirectoryHandle>>>,
}
impl OpfsFilesystem {
pub fn new() -> Self {
Self::default()
}
async fn root_handle(&self) -> Result<FileSystemDirectoryHandle> {
let cached = self.root.borrow().clone();
if let Some(h) = cached {
return Ok(h);
}
let window = web_sys::window()
.ok_or_else(|| Error::fs("getDirectory", "", "no window: not in a browser"))?;
let storage = window.navigator().storage();
let promise = storage.get_directory();
let val = JsFuture::from(promise)
.await
.map_err(|e| Error::fs("getDirectory", "", format!("getDirectory: {}", js_err(&e))))?;
let handle: FileSystemDirectoryHandle = val
.dyn_into()
.map_err(|_| {
Error::fs("getDirectory", "", "getDirectory: not a FileSystemDirectoryHandle")
})?;
*self.root.borrow_mut() = Some(handle.clone());
Ok(handle)
}
async fn resolve_parent(
&self,
path: &str,
create_dirs: bool,
) -> Result<(FileSystemDirectoryHandle, Option<String>)> {
let parts = split_path(path);
if parts.is_empty() {
return Ok((self.root_handle().await?, None));
}
let mut dir = self.root_handle().await?;
for component in &parts[..parts.len() - 1] {
dir = get_subdir(&dir, component, create_dirs).await?;
}
Ok((dir, Some(parts.last().unwrap().clone())))
}
async fn resolve_dir(&self, path: &str) -> Result<FileSystemDirectoryHandle> {
let parts = split_path(path);
let mut dir = self.root_handle().await?;
for component in &parts {
dir = get_subdir(&dir, component, false).await?;
}
Ok(dir)
}
async fn write_atomic_stream(&self, path: &str, bytes: &[u8]) -> Result<()> {
let (parent, name) = self.resolve_parent(path, true).await?;
let name = name.ok_or_else(|| {
Error::fs("write_atomic", path, format!("write_atomic({path}): path is empty"))
})?;
let file_handle = get_file(&parent, &name, true).await?;
let writable_val = JsFuture::from(file_handle.create_writable())
.await
.map_err(|e| {
Error::fs("createWritable", path, format!("createWritable({path}): {}", js_err(&e)))
})?;
let writable: FileSystemWritableFileStream = writable_val
.dyn_into()
.map_err(|_| Error::fs("createWritable", path, "createWritable: not a writable stream"))?;
let array = Uint8Array::from(bytes);
let write_promise = writable
.write_with_buffer_source(&array)
.map_err(|e| Error::fs("write", path, format!("write({path}): {}", js_err(&e))))?;
JsFuture::from(write_promise)
.await
.map_err(|e| Error::fs("write", path, format!("write({path}): {}", js_err(&e))))?;
JsFuture::from(writable.close())
.await
.map_err(|e| Error::fs("close", path, format!("close({path}): {}", js_err(&e))))?;
Ok(())
}
}
async fn bounded(
ms: u32,
op: &'static str,
path: &str,
fut: impl std::future::Future<Output = Result<()>>,
) -> Result<()> {
use futures_util::future::{select, Either};
match select(Box::pin(fut), Box::pin(crate::runtime::sleep_ms(ms))).await {
Either::Left((res, _)) => res,
Either::Right(((), _)) => Err(Error::fs(
op,
path,
format!("{op}({path}): timed out after {ms}ms (engine OPFS stall) — data NOT saved"),
)),
}
}
thread_local! {
static BROKER: RefCell<Option<Broker>> = const { RefCell::new(None) };
static BROKER_DEAD: Cell<bool> = const { Cell::new(false) };
}
struct Broker {
worker: Worker,
pending: Rc<RefCell<HashMap<u32, Function>>>,
next_id: Cell<u32>,
_onmessage: Closure<dyn FnMut(MessageEvent)>,
}
enum BrokerWrite {
Done(Result<()>),
Unsupported,
}
fn engine_prefers_broker() -> bool {
if BROKER_DEAD.with(Cell::get) {
return false;
}
let Some(win) = web_sys::window() else { return false };
if Reflect::get(&win, &JsValue::from_str("LH_FORCE_WORKER_FS"))
.map(|v| v.is_truthy())
.unwrap_or(false)
{
return true;
}
Reflect::get(&win.navigator(), &JsValue::from_str("vendor"))
.ok()
.and_then(|v| v.as_string())
.is_some_and(|v| v.starts_with("Apple"))
}
fn broker_spawn() -> Option<Broker> {
let worker = Worker::new("/opfs-worker.js").ok()?;
let pending: Rc<RefCell<HashMap<u32, Function>>> = Rc::new(RefCell::new(HashMap::new()));
let map = Rc::clone(&pending);
let onmessage = Closure::<dyn FnMut(MessageEvent)>::new(move |e: MessageEvent| {
let data = e.data();
let id = Reflect::get(&data, &JsValue::from_str("id"))
.ok()
.and_then(|v| v.as_f64())
.unwrap_or(-1.0) as u32;
let resolve = map.borrow_mut().remove(&id);
if let Some(resolve) = resolve {
let _ = resolve.call1(&JsValue::NULL, &data);
}
});
worker.set_onmessage(Some(onmessage.as_ref().unchecked_ref()));
Some(Broker { worker, pending, next_id: Cell::new(1), _onmessage: onmessage })
}
async fn broker_write(path: &str, bytes: &[u8]) -> BrokerWrite {
let fut = BROKER.with(|slot| {
let mut slot = slot.borrow_mut();
if slot.is_none() {
*slot = broker_spawn();
}
let b = slot.as_ref()?;
let id = b.next_id.get();
b.next_id.set(id.wrapping_add(1));
let pending = Rc::clone(&b.pending);
let promise = Promise::new(&mut |resolve, _reject| {
pending.borrow_mut().insert(id, resolve);
});
let msg = Object::new();
let _ = Reflect::set(&msg, &JsValue::from_str("id"), &JsValue::from_f64(id as f64));
let _ = Reflect::set(&msg, &JsValue::from_str("path"), &JsValue::from_str(path));
let _ = Reflect::set(&msg, &JsValue::from_str("bytes"), &Uint8Array::from(bytes).buffer());
if b.worker.post_message(&msg).is_err() {
b.pending.borrow_mut().remove(&id);
return None;
}
Some(JsFuture::from(promise))
});
let Some(fut) = fut else {
BROKER_DEAD.with(|d| d.set(true));
return BrokerWrite::Unsupported;
};
use futures_util::future::{select, Either};
let reply = match select(Box::pin(fut), Box::pin(crate::runtime::sleep_ms(WRITE_TIMEOUT_MS)))
.await
{
Either::Left((Ok(v), _)) => v,
Either::Left((Err(e), _)) => {
return BrokerWrite::Done(Err(Error::fs(
"write_atomic",
path,
format!("write broker({path}): {}", js_err(&e)),
)));
}
Either::Right(((), _)) => {
return BrokerWrite::Done(Err(Error::fs(
"write_atomic",
path,
format!("write broker({path}): timed out after {WRITE_TIMEOUT_MS}ms — data NOT saved"),
)));
}
};
let get = |k: &str| Reflect::get(&reply, &JsValue::from_str(k)).ok();
if get("ok").and_then(|v| v.as_bool()) == Some(true) {
return BrokerWrite::Done(Ok(()));
}
if get("unsupported").and_then(|v| v.as_bool()) == Some(true) {
BROKER_DEAD.with(|d| d.set(true));
return BrokerWrite::Unsupported;
}
let err = get("err").and_then(|v| v.as_string()).unwrap_or_else(|| "unknown".into());
BrokerWrite::Done(Err(Error::fs("write_atomic", path, format!("write broker({path}): {err}"))))
}
#[async_trait(?Send)]
impl Filesystem for OpfsFilesystem {
async fn read(&self, path: &str) -> Result<Vec<u8>> {
let (parent, name) = self.resolve_parent(path, false).await?;
let name =
name.ok_or_else(|| Error::fs("read", path, format!("read({path}): path is empty")))?;
let file_handle = get_file(&parent, &name, false).await?;
let file_val = JsFuture::from(file_handle.get_file())
.await
.map_err(|e| Error::fs("getFile", path, format!("getFile({path}): {}", js_err(&e))))?;
let file: File = file_val
.dyn_into()
.map_err(|_| Error::fs("getFile", path, format!("getFile({path}): not a File")))?;
let buf = JsFuture::from(file.array_buffer())
.await
.map_err(|e| {
Error::fs("arrayBuffer", path, format!("arrayBuffer({path}): {}", js_err(&e)))
})?;
let array = Uint8Array::new(&buf);
Ok(array.to_vec())
}
async fn write_atomic(&self, path: &str, bytes: &[u8]) -> Result<()> {
if engine_prefers_broker() {
match broker_write(path, bytes).await {
BrokerWrite::Done(res) => return res,
BrokerWrite::Unsupported => {}
}
}
bounded(WRITE_TIMEOUT_MS, "write_atomic", path, self.write_atomic_stream(path, bytes))
.await
}
async fn metadata(&self, path: &str) -> Result<Option<Metadata>> {
let (parent, name) = self.resolve_parent(path, false).await?;
let Some(name) = name else {
return Ok(Some(Metadata {
kind: EntryKind::Directory,
size: 0,
}));
};
match get_file(&parent, &name, false).await {
Ok(fh) => {
let file_val = JsFuture::from(fh.get_file())
.await
.map_err(|e| Error::fs("getFile", path, format!("getFile({path}): {}", js_err(&e))))?;
let file: File = file_val
.dyn_into()
.map_err(|_| Error::fs("getFile", path, format!("getFile({path}): not a File")))?;
Ok(Some(Metadata {
kind: EntryKind::File,
size: file.size() as u64,
}))
}
Err(_) => match get_subdir(&parent, &name, false).await {
Ok(_) => Ok(Some(Metadata {
kind: EntryKind::Directory,
size: 0,
})),
Err(_) => Ok(None),
},
}
}
async fn read_dir(&self, path: &str) -> Result<Vec<DirEntry>> {
let dir = self.resolve_dir(path).await?;
let mut entries = collect_entries(&dir).await?;
entries.sort_by(|a, b| a.name.cmp(&b.name));
Ok(entries)
}
async fn walk(&self, path: &str, max_depth: Option<usize>) -> Result<Vec<WalkEntry>> {
let root = self.resolve_dir(path).await?;
let mut out = Vec::new();
out.push(WalkEntry {
path: path.trim_end_matches('/').to_string(),
kind: EntryKind::Directory,
size: None,
});
walk_dir(&root, path, 1, max_depth, &mut out).await?;
Ok(out)
}
async fn delete(&self, path: &str) -> Result<()> {
let (parent, name) = self.resolve_parent(path, false).await?;
let name =
name.ok_or_else(|| {
Error::fs("delete", path, format!("delete({path}): cannot delete OPFS root"))
})?;
let opts = FileSystemRemoveOptions::new();
opts.set_recursive(true);
let promise = parent.remove_entry_with_options(&name, &opts);
JsFuture::from(promise)
.await
.map_err(|e| {
Error::fs("removeEntry", path, format!("removeEntry({path}): {}", js_err(&e)))
})?;
Ok(())
}
}
fn split_path(path: &str) -> Vec<String> {
path.split('/')
.filter(|s| !s.is_empty() && *s != ".")
.map(|s| s.to_string())
.collect()
}
async fn get_subdir(
parent: &FileSystemDirectoryHandle,
name: &str,
create: bool,
) -> Result<FileSystemDirectoryHandle> {
let opts = FileSystemGetDirectoryOptions::new();
opts.set_create(create);
let promise = parent.get_directory_handle_with_options(name, &opts);
let val = JsFuture::from(promise)
.await
.map_err(|e| {
Error::fs(
"getDirectoryHandle",
name,
format!("getDirectoryHandle({name}): {}", js_err(&e)),
)
})?;
val.dyn_into().map_err(|_| {
Error::fs("getDirectoryHandle", name, format!("getDirectoryHandle({name}): wrong type"))
})
}
async fn get_file(
parent: &FileSystemDirectoryHandle,
name: &str,
create: bool,
) -> Result<FileSystemFileHandle> {
let opts = FileSystemGetFileOptions::new();
opts.set_create(create);
let promise = parent.get_file_handle_with_options(name, &opts);
let val = JsFuture::from(promise)
.await
.map_err(|e| {
Error::fs("getFileHandle", name, format!("getFileHandle({name}): {}", js_err(&e)))
})?;
val.dyn_into()
.map_err(|_| Error::fs("getFileHandle", name, format!("getFileHandle({name}): wrong type")))
}
async fn collect_entries(dir: &FileSystemDirectoryHandle) -> Result<Vec<DirEntry>> {
let iter_method =
Reflect::get(dir, &JsValue::from_str("entries"))
.map_err(|_| Error::fs("entries", "", "entries"))?;
let iter_fn = iter_method
.dyn_ref::<js_sys::Function>()
.ok_or_else(|| Error::fs("entries", "", "entries() not callable"))?;
let iterator = iter_fn
.call0(dir)
.map_err(|e| Error::fs("entries", "", format!("entries(): {}", js_err(&e))))?;
let next_fn = Reflect::get(&iterator, &JsValue::from_str("next"))
.map_err(|_| Error::fs("iterator.next", "", "iterator.next"))?
.dyn_into::<js_sys::Function>()
.map_err(|_| Error::fs("iterator.next", "", "iterator.next not a function"))?;
let mut out = Vec::new();
loop {
let promise = next_fn
.call0(&iterator)
.map_err(|e| Error::fs("iterator.next", "", format!("iterator.next: {}", js_err(&e))))?;
let result = JsFuture::from(js_sys::Promise::from(promise))
.await
.map_err(|e| {
Error::fs("iterator await", "", format!("iterator await: {}", js_err(&e)))
})?;
let done = Reflect::get(&result, &JsValue::from_str("done"))
.ok()
.and_then(|v| v.as_bool())
.unwrap_or(true);
if done {
break;
}
let value = Reflect::get(&result, &JsValue::from_str("value"))
.map_err(|_| Error::fs("iterator value", "", "iterator value"))?;
let pair: js_sys::Array = value
.dyn_into()
.map_err(|_| Error::fs("entries", "", "entry value not an array"))?;
let name = pair
.get(0)
.as_string()
.ok_or_else(|| Error::fs("entries", "", "entry[0] not a string"))?;
let handle_val = pair.get(1);
let handle: FileSystemHandle = handle_val
.dyn_into()
.map_err(|_| Error::fs("entries", "", "entry[1] not a FileSystemHandle"))?;
let (kind, size) = match handle.kind() {
FileSystemHandleKind::File => {
let fh: FileSystemFileHandle = handle.unchecked_into();
let file_val = JsFuture::from(fh.get_file())
.await
.map_err(|e| Error::fs("getFile", "", format!("getFile: {}", js_err(&e))))?;
let file: File = file_val
.dyn_into()
.map_err(|_| Error::fs("getFile", "", "getFile: not a File"))?;
(EntryKind::File, Some(file.size() as u64))
}
FileSystemHandleKind::Directory => (EntryKind::Directory, None),
_ => (EntryKind::Other, None),
};
out.push(DirEntry { name, kind, size });
}
Ok(out)
}
async fn walk_dir(
dir: &FileSystemDirectoryHandle,
prefix: &str,
depth: usize,
max_depth: Option<usize>,
out: &mut Vec<WalkEntry>,
) -> Result<()> {
if let Some(d) = max_depth {
if depth > d {
return Ok(());
}
}
let entries = collect_entries(dir).await?;
for entry in entries {
if out.len() >= MAX_WALK_ENTRIES {
return Ok(());
}
let path = if prefix.is_empty() || prefix == "/" {
entry.name.clone()
} else {
format!("{}/{}", prefix.trim_end_matches('/'), entry.name)
};
match entry.kind {
EntryKind::File => {
out.push(WalkEntry {
path,
kind: EntryKind::File,
size: entry.size,
});
}
EntryKind::Directory => {
out.push(WalkEntry {
path: path.clone(),
kind: EntryKind::Directory,
size: None,
});
let sub = get_subdir(dir, &entry.name, false).await?;
Box::pin(walk_dir(&sub, &path, depth + 1, max_depth, out)).await?;
}
_ => {
out.push(WalkEntry {
path,
kind: entry.kind,
size: entry.size,
});
}
}
}
Ok(())
}
fn js_err(e: &JsValue) -> String {
if let Some(s) = e.as_string() {
return s;
}
if let Ok(name) = Reflect::get(e, &JsValue::from_str("name")) {
if let Ok(msg) = Reflect::get(e, &JsValue::from_str("message")) {
return format!(
"{}: {}",
name.as_string().unwrap_or_default(),
msg.as_string().unwrap_or_default()
);
}
}
let obj: Object = e.clone().unchecked_into();
obj.to_string().as_string().unwrap_or_else(|| "<js error>".into())
}