feat(leptos): WebSocket singleton with event dispatch
This commit is contained in:
Generated
+1
@@ -498,6 +498,7 @@ dependencies = [
|
||||
"lucide-leptos",
|
||||
"serde",
|
||||
"serde-wasm-bindgen",
|
||||
"serde_json",
|
||||
"shared-types",
|
||||
"wasm-bindgen",
|
||||
"wasm-logger",
|
||||
|
||||
@@ -17,6 +17,7 @@ web-sys = { version = "0.3", features = [
|
||||
"WebSocket",
|
||||
"MessageEvent",
|
||||
"CloseEvent",
|
||||
"ErrorEvent",
|
||||
"CanvasRenderingContext2d",
|
||||
"AudioContext",
|
||||
"AudioBuffer",
|
||||
@@ -44,6 +45,7 @@ web-sys = { version = "0.3", features = [
|
||||
] }
|
||||
gloo-net = "0.6"
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
serde-wasm-bindgen = "0.6"
|
||||
wasm-logger = "0.2"
|
||||
console_error_panic_hook = "0.1"
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
pub mod app;
|
||||
pub mod ui;
|
||||
pub mod ws;
|
||||
|
||||
use wasm_bindgen::prelude::*;
|
||||
|
||||
|
||||
@@ -0,0 +1,140 @@
|
||||
// services/frontend-leptos/frontend/src/ws/context.rs
|
||||
use leptos::prelude::*;
|
||||
use crate::ws::socket::{WsHandle, WsStatus, WsEvent};
|
||||
use shared_types::message::MessageRecord;
|
||||
use shared_types::voice::ActiveSpeaker;
|
||||
use shared_types::media::MediaState;
|
||||
use shared_types::recording::VoiceRecording;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct WsContext {
|
||||
pub handle: std::rc::Rc<WsHandle>,
|
||||
pub status: ReadSignal<WsStatus>,
|
||||
// Per-event callbacks (set externally by feature components)
|
||||
// Wrapped in Rc so cloning shares the same callback slots
|
||||
pub on_message_created: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(MessageRecord)>>>>,
|
||||
pub on_message_updated: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(MessageRecord)>>>>,
|
||||
pub on_message_deleted: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(String)>>>>,
|
||||
pub on_message_analyzed: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(MessageRecord)>>>>,
|
||||
pub on_voice_active_user: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(ActiveSpeaker)>>>>,
|
||||
pub on_voice_recording_uploaded: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(VoiceRecording)>>>>,
|
||||
pub on_media_state: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(MediaState)>>>>,
|
||||
pub on_binary: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(Vec<u8>)>>>>,
|
||||
}
|
||||
|
||||
impl WsContext {
|
||||
pub fn new(url: &str) -> Self {
|
||||
let ws_handle = std::rc::Rc::new(WsHandle::new(url));
|
||||
let status = ws_handle.status;
|
||||
|
||||
let ctx = Self {
|
||||
status,
|
||||
handle: ws_handle,
|
||||
on_message_created: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
on_message_updated: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
on_message_deleted: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
on_message_analyzed: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
on_voice_active_user: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
on_voice_recording_uploaded: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
on_media_state: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
on_binary: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
};
|
||||
|
||||
// Wire up the main event dispatcher
|
||||
let ctx_clone = ctx.clone();
|
||||
ctx.handle.on_event(move |event| {
|
||||
ctx_clone.dispatch_event(event);
|
||||
});
|
||||
|
||||
ctx
|
||||
}
|
||||
|
||||
fn dispatch_event(&self, event: WsEvent) {
|
||||
match event {
|
||||
WsEvent::Text(text) => {
|
||||
// Parse JSON envelope: { type: string, data?: any }
|
||||
if let Ok(parsed) = serde_json::from_str::<serde_json::Value>(&text) {
|
||||
let event_type = parsed["type"].as_str().unwrap_or("").to_string();
|
||||
let data = parsed.get("data");
|
||||
|
||||
match event_type.as_str() {
|
||||
"message_created" => {
|
||||
if let Some(d) = data.and_then(|v| serde_json::from_value::<MessageRecord>(v.clone()).ok()) {
|
||||
if let Some(cb) = self.on_message_created.borrow().as_ref() {
|
||||
cb(d);
|
||||
}
|
||||
}
|
||||
}
|
||||
"message_updated" => {
|
||||
if let Some(d) = data.and_then(|v| serde_json::from_value::<MessageRecord>(v.clone()).ok()) {
|
||||
if let Some(cb) = self.on_message_updated.borrow().as_ref() {
|
||||
cb(d);
|
||||
}
|
||||
}
|
||||
}
|
||||
"message_deleted" => {
|
||||
if let Some(d) = data.and_then(|v| v.as_str().map(String::from)) {
|
||||
if let Some(cb) = self.on_message_deleted.borrow().as_ref() {
|
||||
cb(d);
|
||||
}
|
||||
}
|
||||
}
|
||||
"message_analyzed" => {
|
||||
if let Some(d) = data.and_then(|v| serde_json::from_value::<MessageRecord>(v.clone()).ok()) {
|
||||
if let Some(cb) = self.on_message_analyzed.borrow().as_ref() {
|
||||
cb(d);
|
||||
}
|
||||
}
|
||||
}
|
||||
"voice_active_user" => {
|
||||
if let Some(d) = data.and_then(|v| serde_json::from_value::<ActiveSpeaker>(v.clone()).ok()) {
|
||||
if let Some(cb) = self.on_voice_active_user.borrow().as_ref() {
|
||||
cb(d);
|
||||
}
|
||||
}
|
||||
}
|
||||
"voice_recording_uploaded" => {
|
||||
if let Some(d) = data.and_then(|v| serde_json::from_value::<VoiceRecording>(v.clone()).ok()) {
|
||||
if let Some(cb) = self.on_voice_recording_uploaded.borrow().as_ref() {
|
||||
cb(d);
|
||||
}
|
||||
}
|
||||
}
|
||||
"media_state" => {
|
||||
if let Some(d) = data.and_then(|v| serde_json::from_value::<MediaState>(v.clone()).ok()) {
|
||||
if let Some(cb) = self.on_media_state.borrow().as_ref() {
|
||||
cb(d);
|
||||
}
|
||||
}
|
||||
}
|
||||
_ => {
|
||||
// Unknown event type — log and ignore
|
||||
web_sys::console::log_1(&format!("[WS] unhandled event: {}", event_type).into());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
WsEvent::Binary(data) => {
|
||||
if let Some(cb) = self.on_binary.borrow().as_ref() {
|
||||
cb(data);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn connect(&self) {
|
||||
self.handle.connect();
|
||||
}
|
||||
|
||||
pub fn disconnect(&self) {
|
||||
self.handle.disconnect();
|
||||
}
|
||||
|
||||
pub fn send_text(&self, text: &str) {
|
||||
let _ = self.handle.send_text(text);
|
||||
}
|
||||
|
||||
pub fn send_binary(&self, data: &[u8]) {
|
||||
let _ = self.handle.send_binary(data);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
// services/frontend-leptos/frontend/src/ws/mod.rs
|
||||
pub mod socket;
|
||||
pub mod context;
|
||||
@@ -0,0 +1,152 @@
|
||||
// services/frontend-leptos/frontend/src/ws/socket.rs
|
||||
use leptos::prelude::*;
|
||||
use wasm_bindgen::prelude::*;
|
||||
use wasm_bindgen::JsCast;
|
||||
use web_sys::{WebSocket, MessageEvent, CloseEvent, ErrorEvent};
|
||||
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub enum WsStatus {
|
||||
Disconnected,
|
||||
Connecting,
|
||||
Connected,
|
||||
Error(String),
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum WsEvent {
|
||||
Text(String),
|
||||
Binary(Vec<u8>),
|
||||
}
|
||||
|
||||
pub struct WsHandle {
|
||||
pub status: ReadSignal<WsStatus>,
|
||||
set_status: WriteSignal<WsStatus>,
|
||||
ws: std::cell::RefCell<Option<WebSocket>>,
|
||||
on_event: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(WsEvent)>>>>,
|
||||
url: String,
|
||||
reconnect_attempt: std::cell::Cell<u32>,
|
||||
}
|
||||
|
||||
impl WsHandle {
|
||||
pub fn new(url: &str) -> Self {
|
||||
let (status, set_status) = create_signal(WsStatus::Disconnected);
|
||||
Self {
|
||||
status,
|
||||
set_status,
|
||||
ws: std::cell::RefCell::new(None),
|
||||
on_event: std::rc::Rc::new(std::cell::RefCell::new(None)),
|
||||
url: url.to_string(),
|
||||
reconnect_attempt: std::cell::Cell::new(0),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn on_event<F>(&self, callback: F)
|
||||
where
|
||||
F: Fn(WsEvent) + 'static,
|
||||
{
|
||||
*self.on_event.borrow_mut() = Some(Box::new(callback));
|
||||
}
|
||||
|
||||
pub fn connect(&self) {
|
||||
if self.status.get() == WsStatus::Connected || self.status.get() == WsStatus::Connecting {
|
||||
return;
|
||||
}
|
||||
self.set_status.set(WsStatus::Connecting);
|
||||
|
||||
let url = self.url.clone();
|
||||
let status_clone = self.set_status.clone();
|
||||
// Clone the Rc wrapper (cheap pointer copy, inner RefCell is shared)
|
||||
let event_clone: std::rc::Rc<std::cell::RefCell<Option<Box<dyn Fn(WsEvent)>>>> = self.on_event.clone();
|
||||
let ws_holder = &self.ws as *const std::cell::RefCell<Option<WebSocket>>;
|
||||
let reconnect_attempt = &self.reconnect_attempt as *const std::cell::Cell<u32>;
|
||||
|
||||
// Clone again for closures
|
||||
let status2 = status_clone.clone();
|
||||
let status3 = status_clone.clone();
|
||||
|
||||
match WebSocket::new(&url) {
|
||||
Ok(ws) => {
|
||||
// Store reference
|
||||
unsafe { *(*ws_holder).borrow_mut() = Some(ws.clone()) };
|
||||
|
||||
// onopen
|
||||
let onopen_cb = Closure::<dyn Fn(web_sys::ProgressEvent)>::new(move |_| {
|
||||
status_clone.set(WsStatus::Connected);
|
||||
unsafe { (*reconnect_attempt).set(0) };
|
||||
});
|
||||
ws.set_onopen(Some(onopen_cb.as_ref().unchecked_ref()));
|
||||
onopen_cb.forget();
|
||||
|
||||
// onclose
|
||||
let onclose_cb = Closure::<dyn Fn(CloseEvent)>::new(move |_| {
|
||||
status2.set(WsStatus::Disconnected);
|
||||
unsafe { *(*ws_holder).borrow_mut() = None };
|
||||
});
|
||||
ws.set_onclose(Some(onclose_cb.as_ref().unchecked_ref()));
|
||||
onclose_cb.forget();
|
||||
|
||||
// onerror
|
||||
let onerror_cb = Closure::<dyn Fn(ErrorEvent)>::new(move |e: ErrorEvent| {
|
||||
let msg = e.message();
|
||||
status3.set(WsStatus::Error(msg));
|
||||
});
|
||||
ws.set_onerror(Some(onerror_cb.as_ref().unchecked_ref()));
|
||||
onerror_cb.forget();
|
||||
|
||||
// onmessage
|
||||
let onmsg_cb = Closure::<dyn Fn(MessageEvent)>::new(move |e: MessageEvent| {
|
||||
if let Some(cb) = &*event_clone.borrow() {
|
||||
if let Some(text) = e.data().as_string() {
|
||||
cb(WsEvent::Text(text));
|
||||
} else if let Some(abuf) = e.data().dyn_ref::<js_sys::ArrayBuffer>() {
|
||||
let len = abuf.byte_length() as usize;
|
||||
let u8view = js_sys::Uint8Array::new(abuf);
|
||||
let mut bytes = vec![0u8; len];
|
||||
u8view.copy_to(&mut bytes);
|
||||
cb(WsEvent::Binary(bytes));
|
||||
} else {
|
||||
// Try Blob
|
||||
let data = e.data();
|
||||
let blob = data.dyn_ref::<web_sys::Blob>();
|
||||
if blob.is_some() {
|
||||
// Blob handling would need async FileReader — skip for now
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
ws.set_onmessage(Some(onmsg_cb.as_ref().unchecked_ref()));
|
||||
onmsg_cb.forget();
|
||||
}
|
||||
Err(e) => {
|
||||
status_clone.set(WsStatus::Error(
|
||||
js_sys::Error::from(e).to_string().as_string().unwrap_or_default(),
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub fn disconnect(&self) {
|
||||
if let Some(ws) = self.ws.borrow_mut().take() {
|
||||
ws.close().ok();
|
||||
}
|
||||
self.set_status.set(WsStatus::Disconnected);
|
||||
}
|
||||
|
||||
pub fn send_text(&self, text: &str) -> Result<(), JsValue> {
|
||||
if let Some(ws) = self.ws.borrow().as_ref() {
|
||||
ws.send_with_str(text)
|
||||
} else {
|
||||
Err(JsValue::from_str("WebSocket not connected"))
|
||||
}
|
||||
}
|
||||
|
||||
pub fn send_binary(&self, data: &[u8]) -> Result<(), JsValue> {
|
||||
if let Some(ws) = self.ws.borrow().as_ref() {
|
||||
let array = js_sys::Uint8Array::from(data);
|
||||
let buffer = array.buffer();
|
||||
ws.send_with_array_buffer(&buffer)
|
||||
} else {
|
||||
Err(JsValue::from_str("WebSocket not connected"))
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user