diff --git a/services/frontend-leptos/Cargo.lock b/services/frontend-leptos/Cargo.lock index 6449282..af0a264 100644 --- a/services/frontend-leptos/Cargo.lock +++ b/services/frontend-leptos/Cargo.lock @@ -498,6 +498,7 @@ dependencies = [ "lucide-leptos", "serde", "serde-wasm-bindgen", + "serde_json", "shared-types", "wasm-bindgen", "wasm-logger", diff --git a/services/frontend-leptos/frontend/Cargo.toml b/services/frontend-leptos/frontend/Cargo.toml index c8cabf3..07de2e1 100644 --- a/services/frontend-leptos/frontend/Cargo.toml +++ b/services/frontend-leptos/frontend/Cargo.toml @@ -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" diff --git a/services/frontend-leptos/frontend/src/lib.rs b/services/frontend-leptos/frontend/src/lib.rs index 096f807..ea7fdf6 100644 --- a/services/frontend-leptos/frontend/src/lib.rs +++ b/services/frontend-leptos/frontend/src/lib.rs @@ -1,5 +1,6 @@ pub mod app; pub mod ui; +pub mod ws; use wasm_bindgen::prelude::*; diff --git a/services/frontend-leptos/frontend/src/ws/context.rs b/services/frontend-leptos/frontend/src/ws/context.rs new file mode 100644 index 0000000..e93789f --- /dev/null +++ b/services/frontend-leptos/frontend/src/ws/context.rs @@ -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, + pub status: ReadSignal, + // 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>>>, + pub on_message_updated: std::rc::Rc>>>, + pub on_message_deleted: std::rc::Rc>>>, + pub on_message_analyzed: std::rc::Rc>>>, + pub on_voice_active_user: std::rc::Rc>>>, + pub on_voice_recording_uploaded: std::rc::Rc>>>, + pub on_media_state: std::rc::Rc>>>, + pub on_binary: std::rc::Rc)>>>>, +} + +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::(&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::(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::(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::(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::(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::(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::(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); + } +} diff --git a/services/frontend-leptos/frontend/src/ws/mod.rs b/services/frontend-leptos/frontend/src/ws/mod.rs new file mode 100644 index 0000000..71e685b --- /dev/null +++ b/services/frontend-leptos/frontend/src/ws/mod.rs @@ -0,0 +1,3 @@ +// services/frontend-leptos/frontend/src/ws/mod.rs +pub mod socket; +pub mod context; diff --git a/services/frontend-leptos/frontend/src/ws/socket.rs b/services/frontend-leptos/frontend/src/ws/socket.rs new file mode 100644 index 0000000..1e612e4 --- /dev/null +++ b/services/frontend-leptos/frontend/src/ws/socket.rs @@ -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), +} + +pub struct WsHandle { + pub status: ReadSignal, + set_status: WriteSignal, + ws: std::cell::RefCell>, + on_event: std::rc::Rc>>>, + url: String, + reconnect_attempt: std::cell::Cell, +} + +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(&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>>> = self.on_event.clone(); + let ws_holder = &self.ws as *const std::cell::RefCell>; + let reconnect_attempt = &self.reconnect_attempt as *const std::cell::Cell; + + // 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::::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::::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::::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::::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::() { + 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::(); + 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")) + } + } +}