#![cfg(all(target_arch = "wasm32", feature = "opfs"))]
#![allow(unsafe_code)]
use std::collections::HashMap;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex, Weak};
use futures::channel::oneshot;
use send_wrapper::SendWrapper;
use wasm_bindgen::prelude::*;
use crate::Result;
use crate::errors::PagedbError;
use crate::vfs::traits::Vfs;
use crate::vfs::types::OpenMode;
use super::handle::{OpfsFile, map_err};
use super::lock::{LockMap, OpfsLockHandle};
use super::protocol::{OpfsOp, OpfsRequest, OpfsResponse, OpfsResult};
type Registry = Arc<Mutex<HashMap<u64, oneshot::Sender<OpfsResult>>>>;
struct OpfsVfsInner {
worker: SendWrapper<web_sys::Worker>,
request_registry: Registry,
next_request_id: AtomicU64,
locks: Arc<Mutex<LockMap>>,
weak_self: Mutex<Weak<OpfsVfsInner>>,
_onmessage: SendWrapper<Closure<dyn FnMut(web_sys::MessageEvent)>>,
root: String,
}
pub struct OpfsVfs(Arc<OpfsVfsInner>);
impl Clone for OpfsVfs {
fn clone(&self) -> Self {
OpfsVfs(Arc::clone(&self.0))
}
}
unsafe impl Send for OpfsVfs {}
unsafe impl Sync for OpfsVfs {}
impl OpfsVfs {
pub fn new(worker_url: &str) -> Result<Self> {
Self::with_root(worker_url, "")
}
pub fn with_root(worker_url: &str, root: &str) -> Result<Self> {
let root = super::path::normalize_root(root);
let registry: Registry = Arc::new(Mutex::new(HashMap::new()));
let worker = web_sys::Worker::new(worker_url)
.map_err(|e| PagedbError::Io(std::io::Error::other(format!("{e:?}"))))?;
let registry_cb = Arc::clone(®istry);
let onmessage: Closure<dyn FnMut(web_sys::MessageEvent)> =
Closure::wrap(Box::new(move |event: web_sys::MessageEvent| {
let data = event.data();
match serde_wasm_bindgen::from_value::<OpfsResponse>(data) {
Ok(resp) => {
let sender = {
let mut reg = registry_cb
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
reg.remove(&resp.id)
};
if let Some(tx) = sender {
let _ = tx.send(resp.result);
}
}
Err(e) => {
tracing::warn!("OpfsVfs: failed to deserialize worker response: {:?}", e);
}
}
}));
worker.set_onmessage(Some(onmessage.as_ref().unchecked_ref()));
let inner = Arc::new(OpfsVfsInner {
worker: SendWrapper::new(worker),
request_registry: registry,
next_request_id: AtomicU64::new(1),
locks: Arc::new(Mutex::new(LockMap::default())),
weak_self: Mutex::new(Weak::new()),
_onmessage: SendWrapper::new(onmessage),
root,
});
*inner
.weak_self
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Arc::downgrade(&inner);
Ok(OpfsVfs(inner))
}
fn resolve(&self, path: &str) -> String {
super::path::join_root(&self.0.root, path)
}
pub(crate) async fn dispatch(&self, op: OpfsOp) -> Result<OpfsResult> {
let id = self.0.next_request_id.fetch_add(1, Ordering::Relaxed);
let (tx, rx) = oneshot::channel::<OpfsResult>();
{
let mut reg = self
.0
.request_registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
reg.insert(id, tx);
}
let req = OpfsRequest { id, op };
let js_val = serde_wasm_bindgen::to_value(&req).map_err(|e| {
PagedbError::Io(std::io::Error::other(format!("serialize error: {e:?}")))
})?;
self.0.worker.post_message(&js_val).map_err(|e| {
let mut reg = self
.0
.request_registry
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
reg.remove(&id);
PagedbError::Io(std::io::Error::other(format!("{e:?}")))
})?;
rx.await.map_err(|_| {
PagedbError::Io(std::io::Error::other("worker channel closed unexpectedly"))
})
}
fn arc_inner(&self) -> Option<Arc<OpfsVfsInner>> {
self.0
.weak_self
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.upgrade()
}
}
impl Drop for OpfsVfs {
fn drop(&mut self) {
self.0.worker.terminate();
}
}
impl Vfs for OpfsVfs {
type File = OpfsFile;
type LockHandle = OpfsLockHandle;
async fn open(&self, path: &str, mode: OpenMode) -> Result<Self::File> {
let (create, create_new, read_only) = match mode {
OpenMode::Read => (false, false, true),
OpenMode::ReadWrite => (false, false, false),
OpenMode::CreateNew => (true, true, false),
OpenMode::CreateOrOpen => (true, false, false),
};
let result = self
.dispatch(OpfsOp::Open {
path: self.resolve(path),
create,
create_new,
read_only,
})
.await?;
match result {
OpfsResult::Opened { handle_id } => {
let inner_arc = self.arc_inner().ok_or(PagedbError::Unsupported)?;
Ok(OpfsFile {
handle_id,
vfs: Arc::new(OpfsVfs(inner_arc)),
read_only,
})
}
OpfsResult::Err { reason, kind } => Err(map_err(&reason, kind)),
_ => Err(PagedbError::Unsupported),
}
}
async fn remove(&self, path: &str) -> Result<()> {
match self
.dispatch(OpfsOp::Remove {
path: self.resolve(path),
})
.await?
{
OpfsResult::Ok => Ok(()),
OpfsResult::Err { reason, kind } => Err(map_err(&reason, kind)),
_ => Err(PagedbError::Unsupported),
}
}
async fn rename(&self, from: &str, to: &str) -> Result<()> {
match self
.dispatch(OpfsOp::Rename {
from: self.resolve(from),
to: self.resolve(to),
})
.await?
{
OpfsResult::Ok => Ok(()),
OpfsResult::Err { reason, kind } => Err(map_err(&reason, kind)),
_ => Err(PagedbError::Unsupported),
}
}
async fn list_dir(&self, path: &str) -> Result<Vec<String>> {
match self
.dispatch(OpfsOp::ListDir {
path: self.resolve(path),
})
.await?
{
OpfsResult::Entries { names } => Ok(names),
OpfsResult::Err { reason, kind } => Err(map_err(&reason, kind)),
_ => Err(PagedbError::Unsupported),
}
}
async fn mkdir_all(&self, path: &str) -> Result<()> {
match self
.dispatch(OpfsOp::MkdirAll {
path: self.resolve(path),
})
.await?
{
OpfsResult::Ok => Ok(()),
OpfsResult::Err { reason, kind } => Err(map_err(&reason, kind)),
_ => Err(PagedbError::Unsupported),
}
}
async fn sync_dir(&self, _path: &str) -> Result<()> {
Ok(())
}
async fn lock_exclusive(&self, path: &str) -> Result<Self::LockHandle> {
let resolved = self.resolve(path);
let acquired = self
.0
.locks
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.try_exclusive(&resolved);
if acquired {
Ok(OpfsLockHandle {
path: resolved,
locks: Arc::clone(&self.0.locks),
})
} else {
Err(PagedbError::AlreadyLocked)
}
}
async fn lock_shared(&self, path: &str) -> Result<Self::LockHandle> {
let resolved = self.resolve(path);
let acquired = self
.0
.locks
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.try_shared(&resolved);
if acquired {
Ok(OpfsLockHandle {
path: resolved,
locks: Arc::clone(&self.0.locks),
})
} else {
Err(PagedbError::AlreadyLocked)
}
}
}