mod animation;
mod dom;
mod file;
mod history;
mod timer;
mod websocket;
use std::collections::{HashMap, VecDeque};
use std::io;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use web_sys::console;
use crate::{Event, Interest, RawFd, Reactor};
pub use self::animation::WebAnimationFrame;
pub use self::dom::{
BrowserFiles, CanvasSize, CompositionMetadata, ContentBoxPoint, DropFiles, DropMetadata,
DroppedFile, DroppedFileAccess, ElementSize, KeyboardMetadata, PointerMetadata,
PointerModifiers, PointerType, RgbaFrame, TextInputMetadata, TextSelection,
TextSelectionDirection, WebCanvas, WebClipboard, WebDocument, WebElement, WebEvent,
WebEventListener, WebGpuCanvas, WheelDeltaMode, WheelMetadata,
};
pub use self::file::{MAX_READ_BYTES, WebFile};
pub use self::history::{HistoryListener, MAX_HISTORY_PATH_BYTES, WebHistory};
pub use self::timer::WebTimer;
pub use crate::local_task::LocalTaskHandle;
pub use crate::websocket_state::{WebSocketLimits, WebSocketOpen, WebSocketReceive};
use self::websocket::EVENT_QUEUE_CAPACITY;
use self::websocket::WebSocketConnection;
pub fn spawn_local<F>(future: F)
where
F: std::future::Future<Output = ()> + 'static,
{
wasm_bindgen_futures::spawn_local(future);
}
#[must_use = "retain the handle to cancel the browser task"]
pub fn spawn_local_with_handle<F>(future: F) -> LocalTaskHandle
where
F: std::future::Future<Output = ()> + 'static,
{
let (handle, future) = crate::local_task::cancellable(future);
wasm_bindgen_futures::spawn_local(future);
handle
}
pub struct WebReactor {
pending_events: Arc<Mutex<VecDeque<Event>>>,
websockets: HashMap<RawFd, WebSocketConnection>,
next_fd: RawFd,
fd_interests: Arc<Mutex<HashMap<RawFd, Interest>>>,
}
impl WebReactor {
pub fn new() -> io::Result<Self> {
console::log_1(&"Initializing Moirai WebAssembly reactor".into());
Ok(Self {
pending_events: Arc::new(Mutex::new(VecDeque::with_capacity(EVENT_QUEUE_CAPACITY))),
websockets: HashMap::new(),
next_fd: 1,
fd_interests: Arc::new(Mutex::new(HashMap::new())),
})
}
fn allocate_fd(&mut self) -> RawFd {
let fd = self.next_fd;
self.next_fd += 1;
fd
}
pub fn create_websocket(&mut self, url: &str) -> io::Result<RawFd> {
self.create_websocket_with_limits(url, WebSocketLimits::default())
}
pub fn create_websocket_with_limits(
&mut self,
url: &str,
limits: WebSocketLimits,
) -> io::Result<RawFd> {
let fd = self.allocate_fd();
let websocket = WebSocketConnection::new(
fd,
url,
limits,
Arc::clone(&self.pending_events),
Arc::clone(&self.fd_interests),
)?;
self.websockets.insert(fd, websocket);
Ok(fd)
}
pub fn websocket_send(&self, fd: RawFd, data: &[u8]) -> io::Result<()> {
let websocket = self
.websockets
.get(&fd)
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "WebSocket not found"))?;
websocket.send(data)
}
pub fn websocket_recv(&self, fd: RawFd) -> io::Result<Vec<u8>> {
let websocket = self
.websockets
.get(&fd)
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "WebSocket not found"))?;
websocket.receive()
}
pub fn websocket_recv_async(&self, fd: RawFd) -> io::Result<WebSocketReceive> {
let websocket = self
.websockets
.get(&fd)
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "WebSocket not found"))?;
Ok(websocket.receive_async())
}
pub fn websocket_open_async(&self, fd: RawFd) -> io::Result<WebSocketOpen> {
let websocket = self
.websockets
.get(&fd)
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "WebSocket not found"))?;
Ok(websocket.open_async())
}
pub fn websocket_close(&mut self, fd: RawFd) -> io::Result<()> {
let websocket = self
.websockets
.remove(&fd)
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "WebSocket not found"))?;
self.fd_interests
.lock()
.map_err(|_| io::Error::other("WebSocket interest lock is poisoned"))?
.remove(&fd);
self.pending_events
.lock()
.map_err(|_| io::Error::other("WebSocket event lock is poisoned"))?
.retain(|event| event.fd != fd);
drop(websocket);
Ok(())
}
}
impl Reactor for WebReactor {
fn register_fd(&self, fd: RawFd, interest: Interest) -> io::Result<()> {
if !self.websockets.contains_key(&fd) {
return Err(io::Error::new(
io::ErrorKind::NotFound,
"WebSocket not found",
));
}
self.fd_interests
.lock()
.map_err(|_| io::Error::other("WebSocket interest lock is poisoned"))?
.insert(fd, interest);
Ok(())
}
fn unregister_fd(&self, fd: RawFd) -> io::Result<()> {
self.fd_interests
.lock()
.map_err(|_| io::Error::other("WebSocket interest lock is poisoned"))?
.remove(&fd);
Ok(())
}
fn poll_events(&self, _timeout: Option<Duration>) -> io::Result<Vec<Event>> {
let mut pending_events = self
.pending_events
.lock()
.map_err(|_| io::Error::other("WebSocket event lock is poisoned"))?;
Ok(pending_events.drain(..).collect())
}
fn wake(&self) -> io::Result<()> {
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_web_reactor_creation() {
WebReactor::new().expect("a web reactor must be creatable");
}
}