feat: presence persistence (save/load/remove), OFFLINE on disconnect, voice-node leave signalling
This commit is contained in:
parent
742cf03a5f
commit
447d65df1b
4 changed files with 253 additions and 6 deletions
|
|
@ -3,6 +3,7 @@ pub mod dms;
|
||||||
pub mod e2ee_dms;
|
pub mod e2ee_dms;
|
||||||
pub mod guilds;
|
pub mod guilds;
|
||||||
pub mod messages;
|
pub mod messages;
|
||||||
|
pub mod presence;
|
||||||
pub mod social;
|
pub mod social;
|
||||||
|
|
||||||
use anyhow::Result;
|
use anyhow::Result;
|
||||||
|
|
@ -146,7 +147,15 @@ impl Storage {
|
||||||
target_type TEXT, reason TEXT, changes TEXT,
|
target_type TEXT, reason TEXT, changes TEXT,
|
||||||
created_at INTEGER NOT NULL
|
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)
|
.execute(p)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
@ -253,7 +262,15 @@ impl Storage {
|
||||||
target_type TEXT, reason TEXT, changes TEXT,
|
target_type TEXT, reason TEXT, changes TEXT,
|
||||||
created_at BIGINT NOT NULL
|
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)
|
.execute(p)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
|
||||||
127
gateway/src/domain/storage/presence.rs
Normal file
127
gateway/src/domain/storage/presence.rs
Normal file
|
|
@ -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<String>,
|
||||||
|
pub activity_text: Option<String>,
|
||||||
|
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<Vec<PresenceRow>> {
|
||||||
|
match &self.pool {
|
||||||
|
Pool::Sqlite(p) => {
|
||||||
|
let rows = sqlx::query_as::<_, (String, String, String, Option<String>, Option<String>, 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<String>, Option<String>, 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(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -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(
|
pub async fn handle_presence_update(
|
||||||
_stream: &mut (impl AsyncRead + AsyncWrite + Unpin),
|
_stream: &mut (impl AsyncRead + AsyncWrite + Unpin),
|
||||||
seq: &mut u32,
|
seq: &mut u32,
|
||||||
|
|
@ -27,11 +41,19 @@ pub async fn handle_presence_update(
|
||||||
.await
|
.await
|
||||||
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
|
.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 {
|
let info = PresenceInfo {
|
||||||
user_id: sess.user_id.clone(),
|
user_id: sess.user_id.clone(),
|
||||||
nickname: sess.nickname.clone(),
|
nickname: sess.nickname.clone(),
|
||||||
status: req.status.clone(),
|
status: status.clone(),
|
||||||
activity_type: req.activity_type,
|
activity_type: activity_type.clone(),
|
||||||
activity_text: req.activity_text,
|
activity_text: req.activity_text,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
|
@ -41,7 +63,20 @@ pub async fn handle_presence_update(
|
||||||
.await
|
.await
|
||||||
.insert(sess.user_id.clone(), info.clone());
|
.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 {
|
let event = PresenceEventPayload {
|
||||||
user_id: info.user_id,
|
user_id: info.user_id,
|
||||||
nickname: info.nickname,
|
nickname: info.nickname,
|
||||||
|
|
|
||||||
|
|
@ -14,7 +14,8 @@ use tokio::sync::broadcast;
|
||||||
use tracing::{debug, error, info, warn};
|
use tracing::{debug, error, info, warn};
|
||||||
|
|
||||||
use domain::{channels, config, session, storage};
|
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<config::Config>, voice_member_tx: Option<VoiceMemberTx>) -> Result<()> {
|
pub async fn run(cfg: Arc<config::Config>, voice_member_tx: Option<VoiceMemberTx>) -> Result<()> {
|
||||||
let private_mode = cfg.is_private();
|
let private_mode = cfg.is_private();
|
||||||
|
|
@ -76,6 +77,23 @@ pub async fn run(cfg: Arc<config::Config>, voice_member_tx: Option<VoiceMemberTx
|
||||||
voice_member_tx,
|
voice_member_tx,
|
||||||
);
|
);
|
||||||
|
|
||||||
|
if let Ok(rows) = state.storage.load_all_presences().await {
|
||||||
|
let mut presences = state.presences.write().await;
|
||||||
|
for row in rows {
|
||||||
|
presences.insert(
|
||||||
|
row.user_id.clone(),
|
||||||
|
PresenceInfo {
|
||||||
|
user_id: row.user_id,
|
||||||
|
nickname: row.nickname,
|
||||||
|
status: row.status,
|
||||||
|
activity_type: row.activity_type,
|
||||||
|
activity_text: row.activity_text,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
}
|
||||||
|
info!("loaded {} presences from storage", presences.len());
|
||||||
|
}
|
||||||
|
|
||||||
let admin_bind = cfg
|
let admin_bind = cfg
|
||||||
.gateway
|
.gateway
|
||||||
.admin_bind
|
.admin_bind
|
||||||
|
|
@ -190,12 +208,62 @@ async fn handle<S: AsyncRead + AsyncWrite + Unpin + Send + 'static>(
|
||||||
.await;
|
.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)
|
if let Some(ch) = session::get(&state.sessions, &sid)
|
||||||
.await
|
.await
|
||||||
.and_then(|s| s.channel_id)
|
.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;
|
channels::leave(&state.channels, &ch, &sid).await;
|
||||||
handler::channel::broadcast_leave(&state, &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;
|
session::remove(&state.sessions, &sid).await;
|
||||||
state
|
state
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue