From fa74826eeacc23d92a0ab963fecc96324793c910 Mon Sep 17 00:00:00 2001 From: Tom Date: Sun, 2 Aug 2026 10:52:25 -0700 Subject: [PATCH] Merge commit from fork * fix: check membership before subscribing to servers Signed-off-by: IAmTomahawkx * remove active_server when leaving a server. Signed-off-by: IAmTomahawkx --------- Signed-off-by: IAmTomahawkx --- crates/bonfire/src/events/impl.rs | 1 + crates/bonfire/src/events/state.rs | 15 +++++++++++++++ crates/bonfire/src/websocket.rs | 18 ++++++++++++------ 3 files changed, 28 insertions(+), 6 deletions(-) diff --git a/crates/bonfire/src/events/impl.rs b/crates/bonfire/src/events/impl.rs index 96a0f2ab..84a4cdc4 100644 --- a/crates/bonfire/src/events/impl.rs +++ b/crates/bonfire/src/events/impl.rs @@ -567,6 +567,7 @@ impl State { EventV1::ServerMemberLeave { id, user, .. } => { if user == &self.cache.user_id { self.remove_subscription(id).await; + self.remove_active_server(id).await; if let Some(server) = self.cache.servers.remove(id) { for channel in &server.channels { diff --git a/crates/bonfire/src/events/state.rs b/crates/bonfire/src/events/state.rs index bebc0c3a..c65e8ff0 100644 --- a/crates/bonfire/src/events/state.rs +++ b/crates/bonfire/src/events/state.rs @@ -216,4 +216,19 @@ impl State { subscribed.remove(subscription); } + + // Remove a server from the active server state. + pub async fn remove_active_server(&mut self, server_id: &str) -> Option<()> { + let removed = { + let mut lock = self.active_servers.lock().await; + lock.remove(server_id).is_some() + }; + + if removed { + self.remove_subscription(&format!("{server_id}u")).await; + Some(()) + } else { + None + } + } } diff --git a/crates/bonfire/src/websocket.rs b/crates/bonfire/src/websocket.rs index 426f9f27..381ad906 100644 --- a/crates/bonfire/src/websocket.rs +++ b/crates/bonfire/src/websocket.rs @@ -189,6 +189,7 @@ pub async fn client(db: &'static Database, stream: TcpStream, addr: SocketAddr) read, &write, kill_signal_1_s, + db, ); join!(listener, worker); @@ -420,6 +421,7 @@ async fn worker_with_kill_signal( read: WsReader, write: &Mutex, kill_signal_s: async_channel::Sender<()>, + db: &Database, ) { worker( addr, @@ -431,6 +433,7 @@ async fn worker_with_kill_signal( kill_signal_r, read, write, + db, ) .await; kill_signal_s.send(()).await.ok(); @@ -447,6 +450,7 @@ async fn worker( kill_signal_r: async_channel::Receiver<()>, mut read: WsReader, write: &Mutex, + db: &Database, ) { loop { let t1 = read.try_next().fuse(); @@ -508,13 +512,15 @@ async fn worker( .await; } ClientMessage::Subscribe { server_id } => { - let mut servers = active_servers.lock().await; - let has_item = servers.contains_key(&server_id); - servers.insert(server_id, ()); + if db.fetch_member(&server_id, &user_id).await.is_ok() { + let mut servers = active_servers.lock().await; + let has_item = servers.contains_key(&server_id); + servers.insert(server_id, ()); - if !has_item { - // Poke the listener to adjust subscriptions - topic_signal_s.send(()).await.ok(); + if !has_item { + // Poke the listener to adjust subscriptions + topic_signal_s.send(()).await.ok(); + } } } ClientMessage::Ping { data, responded } => {