indicators

This commit is contained in:
pavel 2026-02-26 22:35:27 +01:00
commit af46da77a7
6 changed files with 216 additions and 26 deletions

View file

@ -6,10 +6,22 @@ use std::collections::HashMap;
use tokio::sync::{RwLock, mpsc};
use uuid::Uuid;
#[derive(Clone)]
pub struct ChatClient {
tx: mpsc::UnboundedSender<ServerEvent>,
is_idle: bool,
}
#[derive(Serialize, Clone)]
pub struct OnlineUser {
pub user_id: Uuid,
pub idle: bool,
}
#[derive(Default)]
pub struct ChatHub {
// user_id -> sender
clients: RwLock<HashMap<Uuid, mpsc::UnboundedSender<ServerEvent>>>,
// user_id -> client
clients: RwLock<HashMap<Uuid, ChatClient>>,
}
#[derive(Serialize, Clone)]
@ -26,6 +38,7 @@ pub enum ServerEvent {
UserPresence {
user_id: Uuid,
online: bool,
idle: bool,
},
}
@ -34,12 +47,13 @@ pub enum ServerEvent {
enum ClientEvent {
// Currently no interactive client events for the general chat WS
Ping,
SetIdleStatus { is_idle: bool },
}
impl ChatHub {
pub async fn add_client(&self, user_id: Uuid, tx: mpsc::UnboundedSender<ServerEvent>) {
let mut clients = self.clients.write().await;
clients.insert(user_id, tx);
clients.insert(user_id, ChatClient { tx, is_idle: false });
}
pub async fn remove_client(&self, user_id: Uuid) {
@ -47,33 +61,54 @@ impl ChatHub {
clients.remove(&user_id);
}
pub async fn get_online_users(&self) -> Vec<Uuid> {
pub async fn get_online_users(&self) -> Vec<OnlineUser> {
let clients = self.clients.read().await;
clients.keys().cloned().collect()
clients
.iter()
.map(|(id, client)| OnlineUser {
user_id: *id,
idle: client.is_idle,
})
.collect()
}
pub async fn broadcast_all(&self, event: ServerEvent) {
let clients = self.clients.read().await;
for tx in clients.values() {
let _ = tx.send(event.clone());
for client in clients.values() {
let _ = client.tx.send(event.clone());
}
}
pub async fn broadcast_to_user(&self, user_id: Uuid, event: ServerEvent) {
let clients = self.clients.read().await;
if let Some(tx) = clients.get(&user_id) {
let _ = tx.send(event);
if let Some(client) = clients.get(&user_id) {
let _ = client.tx.send(event);
}
}
pub async fn broadcast_to_many(&self, user_ids: Vec<Uuid>, event: ServerEvent) {
let clients = self.clients.read().await;
for user_id in user_ids {
if let Some(tx) = clients.get(&user_id) {
let _ = tx.send(event.clone());
if let Some(client) = clients.get(&user_id) {
let _ = client.tx.send(event.clone());
}
}
}
pub async fn set_idle_status(&self, user_id: Uuid, is_idle: bool) {
{
let mut clients = self.clients.write().await;
if let Some(client) = clients.get_mut(&user_id) {
client.is_idle = is_idle;
}
}
self.broadcast_all(ServerEvent::UserPresence {
user_id,
online: true,
idle: is_idle,
})
.await;
}
}
pub async fn handle_socket(state: AppState, socket: WebSocket, user_id: Uuid) {
@ -86,6 +121,7 @@ pub async fn handle_socket(state: AppState, socket: WebSocket, user_id: Uuid) {
.broadcast_all(ServerEvent::UserPresence {
user_id,
online: true,
idle: false,
})
.await;
@ -101,8 +137,16 @@ pub async fn handle_socket(state: AppState, socket: WebSocket, user_id: Uuid) {
});
while let Some(Ok(msg)) = ws_receiver.next().await {
if let Message::Close(_) = msg {
break;
match msg {
Message::Close(_) => break,
Message::Text(text) => {
if let Ok(ClientEvent::SetIdleStatus { is_idle }) =
serde_json::from_str::<ClientEvent>(&text)
{
state.chat.set_idle_status(user_id, is_idle).await;
}
}
_ => {}
}
}
@ -113,6 +157,7 @@ pub async fn handle_socket(state: AppState, socket: WebSocket, user_id: Uuid) {
.broadcast_all(ServerEvent::UserPresence {
user_id,
online: false,
idle: false,
})
.await;
}

View file

@ -19,6 +19,7 @@ struct ClientHandle {
is_sharing_video: bool,
is_sharing_screen: bool,
is_speaking: bool,
is_muted: bool,
tx: mpsc::UnboundedSender<ServerEvent>,
}
@ -29,6 +30,7 @@ pub struct VoiceParticipant {
pub is_sharing_video: bool,
pub is_sharing_screen: bool,
pub is_speaking: bool,
pub is_muted: bool,
}
#[derive(Serialize, Clone)]
@ -56,6 +58,10 @@ enum ServerEvent {
user_id: Uuid,
is_speaking: bool,
},
MuteStatusChanged {
user_id: Uuid,
is_muted: bool,
},
Signal {
from_user_id: Uuid,
kind: String,
@ -87,6 +93,9 @@ enum ClientEvent {
SetSpeakingStatus {
is_speaking: bool,
},
SetMuteStatus {
is_muted: bool,
},
PlaySound {
sound_url: String,
},
@ -106,6 +115,7 @@ impl VoiceHub {
is_sharing_video: handle.is_sharing_video,
is_sharing_screen: handle.is_sharing_screen,
is_speaking: handle.is_speaking,
is_muted: handle.is_muted,
})
.collect()
}
@ -128,6 +138,7 @@ impl VoiceHub {
is_sharing_video: peer.is_sharing_video,
is_sharing_screen: peer.is_sharing_screen,
is_speaking: peer.is_speaking,
is_muted: peer.is_muted,
})
.collect::<Vec<_>>();
@ -138,6 +149,7 @@ impl VoiceHub {
is_sharing_video: false,
is_sharing_screen: false,
is_speaking: false,
is_muted: false,
tx,
},
);
@ -253,6 +265,25 @@ impl VoiceHub {
}
}
pub async fn set_mute_status(&self, room_id: Uuid, user_id: Uuid, is_muted: bool) {
let mut rooms = self.rooms.write().await;
let Some(room) = rooms.get_mut(&room_id) else {
return;
};
if let Some(handle) = room.get_mut(&user_id) {
handle.is_muted = is_muted;
for (peer_id, peer) in room.iter() {
if *peer_id != user_id {
let _ = peer
.tx
.send(ServerEvent::MuteStatusChanged { user_id, is_muted });
}
}
}
}
pub async fn play_sound(&self, room_id: Uuid, user_id: Uuid, sound_url: String) {
let rooms = self.rooms.read().await;
let Some(room) = rooms.get(&room_id) else {
@ -329,6 +360,12 @@ pub async fn handle_socket(state: AppState, socket: WebSocket, room_id: Uuid, us
.set_speaking_status(room_id, user_id, is_speaking)
.await;
}
Ok(ClientEvent::SetMuteStatus { is_muted }) => {
state
.voice
.set_mute_status(room_id, user_id, is_muted)
.await;
}
Ok(ClientEvent::PlaySound { sound_url }) => {
state.voice.play_sound(room_id, user_id, sound_url).await;
}