From 7baa5a52fe7a0d71e3083fb2c399e01b42cdb92b Mon Sep 17 00:00:00 2001 From: Alex Emmet <111742636+Alex-Emmet@users.noreply.github.com> Date: Fri, 14 Nov 2025 11:11:30 +0100 Subject: [PATCH] wewo --- src/calls/call_connection.rs | 5 ++--- src/calls/call_group.rs | 18 ++++++++++-------- src/calls/call_manager.rs | 1 + src/calls/caller.rs | 28 ++++++++++++++-------------- 4 files changed, 27 insertions(+), 25 deletions(-) diff --git a/src/calls/call_connection.rs b/src/calls/call_connection.rs index f360727..532462f 100644 --- a/src/calls/call_connection.rs +++ b/src/calls/call_connection.rs @@ -133,7 +133,7 @@ impl CallConnection { let mut user_info = JsonValue::new_object(); let _ = user_info.insert("state", JsonValue::from("muted")); let _ = user_info.insert("streaming", JsonValue::from(false)); - let _ = users.insert(&caller.user_id.to_string(), user_info); + let _ = users.insert(&caller.user_id.read().await.to_string(), user_info); } } response = response.add_data(DataTypes::about, users); @@ -249,8 +249,7 @@ impl CallConnection { .to_json() .to_string(), ); - group.lock().await.get_member(uid).await; - group.lock().await.remove_member(uid); + group.lock().await.disconnect_member(uid).await; } call_manager::remove_inactive().await; } diff --git a/src/calls/call_group.rs b/src/calls/call_group.rs index cf1d32c..85f0a68 100644 --- a/src/calls/call_group.rs +++ b/src/calls/call_group.rs @@ -21,19 +21,21 @@ impl CallGroup { self.callers .insert(user_id, Arc::new(Caller::new(user_id, tx))); } - pub fn disconnect_member(&mut self, user_id: Uuid) { - if let Some(caller) = self.callers.get(&user_id) { - caller.disconnect(); - } - } - pub fn remove_member(&mut self, user_id: Uuid) { self.callers.remove(&user_id); } + pub fn get_member(&self, user_id: &Uuid) -> Option<&Arc> { + self.callers.get(user_id) + } + pub async fn disconnect_member(&mut self, user_id: Uuid) { + if let Some(caller) = self.callers.get(&user_id) { + caller.disconnect().await; + } + } - pub fn is_empty(&self) -> bool { + pub async fn is_empty(&self) -> bool { for caller in self.callers.values() { - if !caller.is_connected() { + if !caller.is_connected().await { return false; } } diff --git a/src/calls/call_manager.rs b/src/calls/call_manager.rs index ea2c27a..fb7fd79 100644 --- a/src/calls/call_manager.rs +++ b/src/calls/call_manager.rs @@ -46,6 +46,7 @@ pub async fn remove_inactive() { .lock() .await .is_empty() + .await { rem.push(cg.clone()); } diff --git a/src/calls/caller.rs b/src/calls/caller.rs index 8fb0fda..645ad4d 100644 --- a/src/calls/caller.rs +++ b/src/calls/caller.rs @@ -1,15 +1,15 @@ use std::sync::Arc; use futures::lock::Mutex; -use tokio::sync::mpsc::UnboundedSender; +use tokio::sync::{RwLock, mpsc::UnboundedSender}; use tungstenite::Utf8Bytes; use uuid::Uuid; pub struct Caller { - pub user_id: Mutex, + pub user_id: RwLock, pub tx: Mutex>>, - pub user_state: Mutex, - pub streaming: Mutex, + pub user_state: RwLock, + pub streaming: RwLock, } #[derive(Clone)] pub enum CallUserState { @@ -21,26 +21,26 @@ pub enum CallUserState { impl Caller { pub fn new(user_id: Uuid, tx: UnboundedSender) -> Self { Self { - user_id: Mutex::new(user_id), + user_id: RwLock::new(user_id), tx: Mutex::new(Some(tx)), - user_state: Mutex::new(CallUserState::Active), - streaming: Mutex::new(false), + user_state: RwLock::new(CallUserState::Active), + streaming: RwLock::new(false), } } - pub fn send(&self, msg: impl Into) { - if let Some(tx) = &self.tx.lock().await { + pub async fn send(&self, msg: impl Into) { + if let Some(tx) = &*self.tx.lock().await { let _ = tx.send(msg.into()); } } - pub fn disconnect(self: Arc) { - self.user_state = CallUserState::Disconnected; - self.tx = None; + pub async fn disconnect(self: &Arc) { + *self.user_state.write().await = CallUserState::Disconnected; + *self.tx.lock().await = None; } - pub fn is_connected(&self) -> bool { - if let CallUserState::Disconnected = self.user_state { + pub async fn is_connected(&self) -> bool { + if let CallUserState::Disconnected = *self.user_state.read().await { false } else { true