use std::sync::mpsc::{self, Receiver};
use dotzuki_engine::link::{NetworkTransport, TransportError};
use serde::Serialize;
use serde::de::DeserializeOwned;
use wasm_bindgen::JsCast;
use wasm_bindgen::closure::Closure;
use wasm_bindgen::prelude::*;
use super::envelope::{Frame, decode_line, encode_line};
pub struct BroadcastChannelTransport<M> {
channel: web_sys::BroadcastChannel,
tag: String,
rx: Receiver<M>,
listener: Option<Closure<dyn FnMut(web_sys::MessageEvent)>>,
}
impl<M> BroadcastChannelTransport<M>
where
M: Serialize + DeserializeOwned + Send + 'static,
{
pub fn new(channel_name: &str) -> Result<Self, TransportError> {
let channel = web_sys::BroadcastChannel::new(channel_name).map_err(|e| {
TransportError::IoError(format!(
"BroadcastChannel '{}' failed: {:?}",
channel_name, e
))
})?;
let tag = random_tag();
let (tx, rx) = mpsc::channel::<M>();
let listener_tag = tag.clone();
let listener_tx = tx.clone();
let listener = Closure::wrap(Box::new(move |event: web_sys::MessageEvent| {
let Some(line) = event.data().as_string() else {
return; };
match decode_line::<Frame<M>>(&line) {
Ok(frame) if frame.is_self(&listener_tag) => {
}
Ok(frame) => {
if listener_tx.send(frame.msg).is_err() {
}
}
Err(e) => {
log::warn!("[link] dropping malformed broadcast frame: {}", e);
}
}
}) as Box<dyn FnMut(_)>);
channel.set_onmessage(Some(listener.as_ref().unchecked_ref()));
Ok(BroadcastChannelTransport {
channel,
tag,
rx,
listener: Some(listener),
})
}
pub fn tag(&self) -> &str {
&self.tag
}
}
impl<M> NetworkTransport<M> for BroadcastChannelTransport<M>
where
M: Serialize + DeserializeOwned + Send + 'static,
{
fn send(&mut self, msg: M) -> Result<(), TransportError> {
let frame = Frame {
from: self.tag.clone(),
msg,
};
let line = encode_line(&frame)?;
self.channel.post_message(&JsValue::from_str(&line)).map_err(|e| {
TransportError::IoError(format!("BroadcastChannel post failed: {:?}", e))
})
}
fn recv(&mut self) -> Result<M, TransportError> {
self.rx.recv().map_err(|_| TransportError::Disconnected)
}
fn try_recv(&mut self) -> Result<Option<M>, TransportError> {
match self.rx.try_recv() {
Ok(msg) => Ok(Some(msg)),
Err(mpsc::TryRecvError::Empty) => Ok(None),
Err(mpsc::TryRecvError::Disconnected) => Err(TransportError::Disconnected),
}
}
}
impl<M> Drop for BroadcastChannelTransport<M> {
fn drop(&mut self) {
self.channel.close();
self.listener.take();
}
}
fn random_tag() -> String {
format!("{:x}", (js_sys::Math::random() * 9_007_199_254_740_992.0) as u64)
}