use crate::actor::{Actor, ActorContext, Addr};
use crate::message::Message;
use async_trait::async_trait;
use std::sync::{Arc, Mutex};
use wasm_bindgen::closure::Closure;
use wasm_bindgen::{JsCast, JsValue};
use web_sys::{MessageEvent, WebSocket};
pub struct WasmWsConn {
ws: WebSocket,
outbox: Arc<Mutex<Vec<String>>>,
}
impl WasmWsConn {
pub fn new(url: &str, ctx: &ActorContext, allow_public_space: bool) -> Self {
let ws = WebSocket::new(url).expect("Failed to create WebSocket");
let peer_id = ctx.peer_id.read().clone();
let router = ctx.router.read().clone();
let addr = ctx.addr.clone();
let stop_ctx = ctx.clone();
let outbox = Arc::new(Mutex::new(Vec::<String>::new()));
let outbox_for_open = outbox.clone();
let peer_id_for_open = peer_id.clone();
let ws_for_open = ws.clone();
let onopen: Closure<dyn FnMut(JsValue)> = Closure::new(move |_event: JsValue| {
let hi = Message::Hi {
from: Addr::noop(),
peer_id: peer_id_for_open.clone(),
is_ack: None,
msg_id: crate::utils::random_string(8),
};
let _ = ws_for_open.send_with_str(&hi.to_string());
let mut queue = outbox_for_open.lock().unwrap();
for text in queue.drain(..) {
let _ = ws_for_open.send_with_str(&text);
}
});
let router_for_msg = router.clone();
let addr_for_msg = addr.clone();
let onmessage: Closure<dyn FnMut(MessageEvent)> =
Closure::new(move |event: MessageEvent| {
if let Some(text) = event.data().as_string() {
match Message::try_from(&text, addr_for_msg.clone(), allow_public_space) {
Ok(msgs) => {
for msg in msgs {
let _ = router_for_msg.send(msg);
}
}
Err(e) => {
web_sys::console::error_1(
&format!("WasmWsConn: parse error: {}", e).into(),
);
}
}
}
});
let onerror: Closure<dyn FnMut(JsValue)> = Closure::new(move |_event: JsValue| {
stop_ctx.stop();
});
let onclose: Closure<dyn FnMut(JsValue)> = Closure::new(move |_event: JsValue| {});
ws.set_onopen(Some(onopen.as_ref().unchecked_ref()));
ws.set_onmessage(Some(onmessage.as_ref().unchecked_ref()));
ws.set_onerror(Some(onerror.as_ref().unchecked_ref()));
ws.set_onclose(Some(onclose.as_ref().unchecked_ref()));
onopen.forget();
onmessage.forget();
onerror.forget();
onclose.forget();
Self { ws, outbox }
}
}
#[async_trait]
impl Actor for WasmWsConn {
async fn pre_start(&mut self, ctx: &ActorContext) {
let _ = ctx.router.read().send(Message::Hi {
from: ctx.addr.clone(),
peer_id: ctx.peer_id.read().clone(),
is_ack: None,
msg_id: crate::utils::random_string(8),
});
}
async fn handle(&mut self, msg: Arc<Message>, _ctx: &ActorContext) {
let mut buf = Vec::with_capacity(256);
msg.to_writer(&mut buf);
let text = std::str::from_utf8(&buf).unwrap_or("");
match self.ws.send_with_str(text) {
Ok(()) => {}
Err(_) => {
self.outbox.lock().unwrap().push(text.to_string());
}
}
}
async fn stopping(&mut self, _ctx: &ActorContext) {
let _ = self.ws.close();
}
}