From 447d65df1b0461584ba14dff845426088cfe56b7 Mon Sep 17 00:00:00 2001 From: loki5512344 Date: Sat, 18 Jul 2026 22:19:48 +0200 Subject: [PATCH] feat: presence persistence (save/load/remove), OFFLINE on disconnect, voice-node leave signalling --- gateway/src/domain/storage/mod.rs | 21 +++- gateway/src/domain/storage/presence.rs | 127 ++++++++++++++++++++++++ gateway/src/handler/content/presence.rs | 41 +++++++- gateway/src/lib.rs | 70 ++++++++++++- 4 files changed, 253 insertions(+), 6 deletions(-) create mode 100644 gateway/src/domain/storage/presence.rs diff --git a/gateway/src/domain/storage/mod.rs b/gateway/src/domain/storage/mod.rs index ba5def5..593eec7 100644 --- a/gateway/src/domain/storage/mod.rs +++ b/gateway/src/domain/storage/mod.rs @@ -3,6 +3,7 @@ pub mod dms; pub mod e2ee_dms; pub mod guilds; pub mod messages; +pub mod presence; pub mod social; use anyhow::Result; @@ -146,7 +147,15 @@ impl Storage { target_type TEXT, reason TEXT, changes TEXT, created_at INTEGER NOT NULL ); - CREATE INDEX IF NOT EXISTS idx_audit_guild ON audit_logs(guild_id, created_at);", + CREATE INDEX IF NOT EXISTS idx_audit_guild ON audit_logs(guild_id, created_at); + CREATE TABLE IF NOT EXISTS presences ( + user_id TEXT PRIMARY KEY, + nickname TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'OFFLINE', + activity_type TEXT, + activity_text TEXT, + last_seen INTEGER NOT NULL + );", ) .execute(p) .await?; @@ -253,7 +262,15 @@ impl Storage { target_type TEXT, reason TEXT, changes TEXT, created_at BIGINT NOT NULL ); - CREATE INDEX IF NOT EXISTS idx_audit_guild ON audit_logs(guild_id, created_at);", + CREATE INDEX IF NOT EXISTS idx_audit_guild ON audit_logs(guild_id, created_at); + CREATE TABLE IF NOT EXISTS presences ( + user_id TEXT PRIMARY KEY, + nickname TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'OFFLINE', + activity_type TEXT, + activity_text TEXT, + last_seen BIGINT NOT NULL + );", ) .execute(p) .await?; diff --git a/gateway/src/domain/storage/presence.rs b/gateway/src/domain/storage/presence.rs new file mode 100644 index 0000000..7da28b6 --- /dev/null +++ b/gateway/src/domain/storage/presence.rs @@ -0,0 +1,127 @@ +use anyhow::Result; + +use super::{Pool, now_ms}; + +pub struct PresenceRow { + pub user_id: String, + pub nickname: String, + pub status: String, + pub activity_type: Option, + pub activity_text: Option, + pub last_seen: i64, +} + +impl super::Storage { + pub async fn save_presence( + &self, + user_id: &str, + nickname: &str, + status: &str, + activity_type: Option<&str>, + activity_text: Option<&str>, + ) -> Result<()> { + let now = now_ms(); + match &self.pool { + Pool::Sqlite(p) => { + sqlx::query( + "INSERT OR REPLACE INTO presences (user_id,nickname,status,activity_type,activity_text,last_seen) \ + VALUES (?,?,?,?,?,?)", + ) + .bind(user_id) + .bind(nickname) + .bind(status) + .bind(activity_type) + .bind(activity_text) + .bind(now) + .execute(p) + .await?; + } + Pool::Postgres(p) => { + sqlx::query( + "INSERT INTO presences (user_id,nickname,status,activity_type,activity_text,last_seen) \ + VALUES ($1,$2,$3,$4,$5,$6) \ + ON CONFLICT (user_id) DO UPDATE SET \ + nickname=EXCLUDED.nickname, status=EXCLUDED.status, \ + activity_type=EXCLUDED.activity_type, activity_text=EXCLUDED.activity_text, \ + last_seen=EXCLUDED.last_seen", + ) + .bind(user_id) + .bind(nickname) + .bind(status) + .bind(activity_type) + .bind(activity_text) + .bind(now) + .execute(p) + .await?; + } + } + Ok(()) + } + + pub async fn load_all_presences(&self) -> Result> { + match &self.pool { + Pool::Sqlite(p) => { + let rows = sqlx::query_as::<_, (String, String, String, Option, Option, i64)>( + "SELECT user_id,nickname,status,activity_type,activity_text,last_seen FROM presences", + ) + .fetch_all(p) + .await?; + Ok(rows + .into_iter() + .map( + |(user_id, nickname, status, activity_type, activity_text, last_seen)| { + PresenceRow { + user_id, + nickname, + status, + activity_type, + activity_text, + last_seen, + } + }, + ) + .collect()) + } + Pool::Postgres(p) => { + let rows = sqlx::query_as::<_, (String, String, String, Option, Option, i64)>( + "SELECT user_id,nickname,status,activity_type,activity_text,last_seen FROM presences", + ) + .fetch_all(p) + .await?; + Ok(rows + .into_iter() + .map( + |(user_id, nickname, status, activity_type, activity_text, last_seen)| { + PresenceRow { + user_id, + nickname, + status, + activity_type, + activity_text, + last_seen, + } + }, + ) + .collect()) + } + } + } + + pub async fn remove_presence(&self, user_id: &str) -> Result<()> { + match &self.pool { + Pool::Sqlite(p) => { + sqlx::query("DELETE FROM presences WHERE user_id=?") + .bind(user_id) + .execute(p) + .await?; + } + Pool::Postgres(p) => { + sqlx::query("DELETE FROM presences WHERE user_id=$1") + .bind(user_id) + .execute(p) + .await?; + } + } + Ok(()) + } +} diff --git a/gateway/src/handler/content/presence.rs b/gateway/src/handler/content/presence.rs index 4af5c5b..2029a69 100644 --- a/gateway/src/handler/content/presence.rs +++ b/gateway/src/handler/content/presence.rs @@ -14,6 +14,20 @@ use crate::{ }, }; +fn valid_status(s: &str) -> bool { + matches!( + s, + "ONLINE" | "IDLE" | "DO_NOT_DISTURB" | "OFFLINE" | "INVISIBLE" + ) +} + +fn valid_activity_type(s: &str) -> bool { + matches!( + s, + "PLAYING" | "LISTENING" | "WATCHING" | "STREAMING" | "CUSTOM" + ) +} + pub async fn handle_presence_update( _stream: &mut (impl AsyncRead + AsyncWrite + Unpin), seq: &mut u32, @@ -27,11 +41,19 @@ pub async fn handle_presence_update( .await .ok_or_else(|| anyhow::anyhow!("session not found"))?; + let status = if valid_status(&req.status) { + req.status.clone() + } else { + "ONLINE".to_string() + }; + + let activity_type = req.activity_type.filter(|a| valid_activity_type(a)); + let info = PresenceInfo { user_id: sess.user_id.clone(), nickname: sess.nickname.clone(), - status: req.status.clone(), - activity_type: req.activity_type, + status: status.clone(), + activity_type: activity_type.clone(), activity_text: req.activity_text, }; @@ -41,7 +63,20 @@ pub async fn handle_presence_update( .await .insert(sess.user_id.clone(), info.clone()); - // Broadcast to everyone (friends/guild-mates will filter client-side for now) + if let Err(e) = state + .storage + .save_presence( + &sess.user_id, + &sess.nickname, + &status, + activity_type.as_deref(), + info.activity_text.as_deref(), + ) + .await + { + tracing::warn!("failed to save presence: {e}"); + } + let event = PresenceEventPayload { user_id: info.user_id, nickname: info.nickname, diff --git a/gateway/src/lib.rs b/gateway/src/lib.rs index 28f3393..0ae185d 100644 --- a/gateway/src/lib.rs +++ b/gateway/src/lib.rs @@ -14,7 +14,8 @@ use tokio::sync::broadcast; use tracing::{debug, error, info, warn}; use domain::{channels, config, session, storage}; -use net::state::{State, VoiceMemberTx}; +use net::state::{BroadcastMsg, State, VoiceMemberTx}; +use proto::{PacketId, PresenceEventPayload, PresenceInfo, encode_packet, to_payload}; pub async fn run(cfg: Arc, voice_member_tx: Option) -> Result<()> { let private_mode = cfg.is_private(); @@ -76,6 +77,23 @@ pub async fn run(cfg: Arc, voice_member_tx: Option( .await; } + let user_id = session::get(&state.sessions, &sid) + .await + .map(|s| s.user_id.clone()); + + if let Some(ref uid) = user_id { + state.presences.write().await.insert( + uid.clone(), + PresenceInfo { + user_id: uid.clone(), + nickname: String::new(), + status: "OFFLINE".into(), + activity_type: None, + activity_text: None, + }, + ); + let _ = state + .storage + .save_presence(uid, "", "OFFLINE", None, None) + .await; + let event = PresenceEventPayload { + user_id: uid.clone(), + nickname: String::new(), + status: "OFFLINE".into(), + activity_type: None, + activity_text: None, + }; + let _ = state.broadcast.send(BroadcastMsg { + channel_id: None, + exclude_session: Some(sid.clone()), + target_session_id: None, + data: encode_packet(PacketId::PresenceEvent, 0, &to_payload(&event)), + }); + } + if let Some(ch) = session::get(&state.sessions, &sid) .await .and_then(|s| s.channel_id) { + let is_voice = channels::get_channel(&state.channels, &ch) + .await + .is_some_and(|c| c.kind == channels::ChannelKind::Voice); + channels::leave(&state.channels, &ch, &sid).await; handler::channel::broadcast_leave(&state, &ch, &sid).await; + + if is_voice + && let Some(tx) = &state.voice_member_tx + && let Some(ref uid) = user_id + { + let event = serde_json::json!({ + "type": "leave", + "channel_id": ch, + "user_id": uid, + }); + let _ = tx.send(event.to_string()); + } } session::remove(&state.sessions, &sid).await; state