diff --git a/gateway/src/handler/channel/create.rs b/gateway/src/handler/channel/create.rs index 755f1a9..b4365c2 100644 --- a/gateway/src/handler/channel/create.rs +++ b/gateway/src/handler/channel/create.rs @@ -99,6 +99,7 @@ pub async fn handle_channel_create( let created = channels::create(&state.channels, &channel_id, &channel_name, kind.clone()).await; // Persist to DB (best-effort, log and continue on failure). + #[allow(clippy::collapsible_if)] if created { if let Err(e) = state .storage @@ -204,6 +205,7 @@ pub async fn handle_channel_delete( let existed = channels::delete(&state.channels, &req.channel_id).await; // Remove from DB (best-effort). + #[allow(clippy::collapsible_if)] if existed { if let Err(e) = state.storage.delete_channel(&req.channel_id).await { warn!("failed to remove channel from storage: {e}"); diff --git a/gateway/src/handler/channel/join.rs b/gateway/src/handler/channel/join.rs index 481bb1f..14401be 100644 --- a/gateway/src/handler/channel/join.rs +++ b/gateway/src/handler/channel/join.rs @@ -110,16 +110,16 @@ pub async fn join( }); } - if let Some(tx) = &state.voice_member_tx { - if let Some(sess) = session::get(&state.sessions, session_id).await { - let event = serde_json::json!({ - "type": "joined", - "channel_id": channel_id, - "session_id": session_id, - "user_id": sess.user_id, - }); - let _ = tx.send(event.to_string()); - } + if let Some(tx) = &state.voice_member_tx + && let Some(sess) = session::get(&state.sessions, session_id).await + { + let event = serde_json::json!({ + "type": "joined", + "channel_id": channel_id, + "session_id": session_id, + "user_id": sess.user_id, + }); + let _ = tx.send(event.to_string()); } info!("session {} joined {channel_id}", &session_id[..8]); diff --git a/gateway/src/handler/channel/leave.rs b/gateway/src/handler/channel/leave.rs index 51521e6..89bf8bc 100644 --- a/gateway/src/handler/channel/leave.rs +++ b/gateway/src/handler/channel/leave.rs @@ -2,7 +2,11 @@ use anyhow::Result; use tokio::io::{AsyncRead, AsyncWrite}; use tracing::info; -use crate::{domain::{channels, session}, net::state::State, proto::SessionCrypto}; +use crate::{ + domain::{channels, session}, + net::state::State, + proto::SessionCrypto, +}; use super::{broadcast_leave, set_channel}; @@ -22,16 +26,16 @@ pub async fn leave( set_channel(state, session_id, None).await; broadcast_leave(state, channel_id, session_id).await; - if let Some(tx) = &state.voice_member_tx { - if let Some(ref uid) = user_id { - let event = serde_json::json!({ - "type": "left", - "channel_id": channel_id, - "session_id": session_id, - "user_id": uid, - }); - let _ = tx.send(event.to_string()); - } + if let Some(tx) = &state.voice_member_tx + && let Some(ref uid) = user_id + { + let event = serde_json::json!({ + "type": "left", + "channel_id": channel_id, + "session_id": session_id, + "user_id": uid, + }); + let _ = tx.send(event.to_string()); } info!("session {} left {channel_id}", &session_id[..8]); diff --git a/gateway/src/handler/dispatch.rs b/gateway/src/handler/dispatch.rs index 6ed6100..e025020 100644 --- a/gateway/src/handler/dispatch.rs +++ b/gateway/src/handler/dispatch.rs @@ -36,36 +36,40 @@ pub async fn dispatch( } PacketId::JoinChannel => { let m = JoinChannelPayload::decode(payload)?; - channel::join(ctx.stream, ctx.seq, session_id, &m.channel_id, ctx.crypto, ctx.state) - .await?; - } - PacketId::LeaveChannel => { - let m = LeaveChannelPayload::decode(payload)?; - channel::leave(ctx.stream, ctx.seq, session_id, &m.channel_id, ctx.crypto, ctx.state) - .await?; - } - PacketId::ChannelCreate => { - channel::handle_channel_create( + channel::join( ctx.stream, ctx.seq, session_id, - payload, + &m.channel_id, ctx.crypto, ctx.state, ) .await?; } - PacketId::ChannelDelete => { - channel::handle_channel_delete( + PacketId::LeaveChannel => { + let m = LeaveChannelPayload::decode(payload)?; + channel::leave( ctx.stream, ctx.seq, session_id, - payload, + &m.channel_id, ctx.crypto, ctx.state, ) .await?; } + PacketId::ChannelCreate => { + channel::handle_channel_create( + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, + ) + .await?; + } + PacketId::ChannelDelete => { + channel::handle_channel_delete( + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, + ) + .await?; + } PacketId::ChannelList => { channel::handle_channel_list(ctx.stream, ctx.seq, session_id, ctx.crypto, ctx.state) .await?; @@ -76,67 +80,37 @@ pub async fn dispatch( } PacketId::DmStart => { direct_message::handle_dm_start( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::DmMessage => { direct_message::handle_dm_message( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::DmHistory => { direct_message::handle_dm_history( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::DmReadAck => { direct_message::handle_dm_read_ack( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildCreate => { guild::handle_guild_create( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildDelete => { guild::handle_guild_delete( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } @@ -146,209 +120,115 @@ pub async fn dispatch( } PacketId::GuildMemberJoin => { guild::handle_guild_member_join( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildMemberLeave => { guild::handle_guild_member_leave( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildMemberKick => { guild::handle_guild_member_kick( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::RoleCreate => { guild::handle_role_create( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::RoleDelete => { guild::handle_role_delete( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::InviteCreate => { guild::handle_invite_create( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::InviteAccept => { guild::handle_invite_accept( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::InviteDelete => { guild::handle_invite_delete( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildAuditLogFetch => { guild::handle_audit_log_fetch( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildMemberListFetch => { guild::handle_member_list_fetch( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildRoleAssign => { guild::handle_role_assign( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildRoleUnassign => { guild::handle_role_unassign( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::GuildRoleListFetch => { guild::handle_role_list_fetch( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::PresenceUpdate => { content::presence::handle_presence_update( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::PresenceSync => { content::presence::handle_presence_sync( - ctx.stream, - ctx.seq, - session_id, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, ctx.crypto, ctx.state, ) .await?; } PacketId::FriendRequest => { friends::handle_friend_request( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::FriendAccept => { friends::handle_friend_accept( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::FriendDecline => { friends::handle_friend_decline( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::FriendRemove => { friends::handle_friend_remove( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } @@ -358,23 +238,13 @@ pub async fn dispatch( } PacketId::BlockUser => { friends::handle_block_user( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } PacketId::UnblockUser => { friends::handle_unblock_user( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } @@ -401,12 +271,7 @@ pub async fn dispatch( } PacketId::ReadReceipt => { content::handle_read_receipt( - ctx.stream, - ctx.seq, - session_id, - payload, - ctx.crypto, - ctx.state, + ctx.stream, ctx.seq, session_id, payload, ctx.crypto, ctx.state, ) .await?; } diff --git a/gateway/src/handler/friends/decline.rs b/gateway/src/handler/friends/decline.rs index b16356e..102360c 100644 --- a/gateway/src/handler/friends/decline.rs +++ b/gateway/src/handler/friends/decline.rs @@ -30,7 +30,11 @@ pub async fn handle_friend_decline( stream, PacketId::FriendDecline, seq, - &to_payload(&SimpleResponsePayload { ok: true, user_id: None, role_id: None }), + &to_payload(&SimpleResponsePayload { + ok: true, + user_id: None, + role_id: None, + }), crypto, ) .await?; diff --git a/gateway/src/handler/friends/manage.rs b/gateway/src/handler/friends/manage.rs index 7302502..9ae4c17 100644 --- a/gateway/src/handler/friends/manage.rs +++ b/gateway/src/handler/friends/manage.rs @@ -33,7 +33,11 @@ pub async fn handle_friend_remove( stream, PacketId::FriendRemove, seq, - &to_payload(&SimpleResponsePayload { ok: true, user_id: None, role_id: None }), + &to_payload(&SimpleResponsePayload { + ok: true, + user_id: None, + role_id: None, + }), crypto, ) .await?; @@ -62,7 +66,11 @@ pub async fn handle_block_user( stream, PacketId::BlockUser, seq, - &to_payload(&SimpleResponsePayload { ok: true, user_id: None, role_id: None }), + &to_payload(&SimpleResponsePayload { + ok: true, + user_id: None, + role_id: None, + }), crypto, ) .await?; @@ -91,7 +99,11 @@ pub async fn handle_unblock_user( stream, PacketId::UnblockUser, seq, - &to_payload(&SimpleResponsePayload { ok: true, user_id: None, role_id: None }), + &to_payload(&SimpleResponsePayload { + ok: true, + user_id: None, + role_id: None, + }), crypto, ) .await?; diff --git a/gateway/src/handler/guild/list.rs b/gateway/src/handler/guild/list.rs index 36b2e6b..f36d79b 100644 --- a/gateway/src/handler/guild/list.rs +++ b/gateway/src/handler/guild/list.rs @@ -4,7 +4,9 @@ use tokio::io::{AsyncRead, AsyncWrite}; use crate::{ domain::session, net::{io, state::State}, - proto::{GuildInfo, GuildListPayload, PacketId, SessionCrypto, UserRoleUpdatePayload, to_payload}, + proto::{ + GuildInfo, GuildListPayload, PacketId, SessionCrypto, UserRoleUpdatePayload, to_payload, + }, }; pub async fn handle_guild_list( diff --git a/gateway/src/net/handshake.rs b/gateway/src/net/handshake.rs index 70282e0..789f886 100644 --- a/gateway/src/net/handshake.rs +++ b/gateway/src/net/handshake.rs @@ -83,7 +83,11 @@ pub async fn run( return Err(anyhow::anyhow!("auth failed")); } - if state.storage.is_banned(&hex::encode(&msg.client_pubkey)).await? { + if state + .storage + .is_banned(&hex::encode(&msg.client_pubkey)) + .await? + { state.metrics.inc(&state.metrics.auth_failures); io::send_error(stream, seq, proto::ErrorCode::AuthFailed, "banned").await?; return Err(anyhow::anyhow!("banned")); diff --git a/gateway/src/net/state.rs b/gateway/src/net/state.rs index 282faf2..a93ad98 100644 --- a/gateway/src/net/state.rs +++ b/gateway/src/net/state.rs @@ -1,7 +1,7 @@ use std::collections::HashMap; -use std::sync::atomic::AtomicUsize; use std::sync::Arc; -use tokio::sync::{broadcast, RwLock}; +use std::sync::atomic::AtomicUsize; +use tokio::sync::{RwLock, broadcast}; use crate::admin::metrics::Metrics; use crate::bootstrap::server_identity::ServerIdentity; diff --git a/gateway/src/proto/crypto/cipher.rs b/gateway/src/proto/crypto/cipher.rs index e02dfa8..8944105 100644 --- a/gateway/src/proto/crypto/cipher.rs +++ b/gateway/src/proto/crypto/cipher.rs @@ -1,6 +1,6 @@ use chacha20poly1305::{ - aead::{Aead, Payload}, ChaCha20Poly1305, Key, KeyInit, Nonce, + aead::{Aead, Payload}, }; pub(super) fn make_nonce(cid: &[u8; 8], seq: u64) -> [u8; 12] { diff --git a/gateway/src/proto/mod.rs b/gateway/src/proto/mod.rs index d8ad057..f0c1e52 100644 --- a/gateway/src/proto/mod.rs +++ b/gateway/src/proto/mod.rs @@ -12,7 +12,7 @@ pub use crypto::SessionCrypto; // Re-export everything so `crate::proto::X` works as before pub use framing::{encode_packet, to_payload}; -pub use packet::{flags, ErrorCode, PacketHeader, PacketId}; +pub use packet::{ErrorCode, PacketHeader, PacketId, flags}; // ─── Backward-compat type aliases ──────────────────────────────── // Handshake diff --git a/serverd/src/main.rs b/serverd/src/main.rs index 92d6889..8ecfa29 100644 --- a/serverd/src/main.rs +++ b/serverd/src/main.rs @@ -1,12 +1,8 @@ -mod voice_membership; - use anyhow::Result; use std::sync::Arc; use tokio::sync::broadcast; use tracing::info; - - #[tokio::main] async fn main() -> Result<()> { tracing_subscriber::fmt() @@ -22,10 +18,7 @@ async fn main() -> Result<()> { let cfg: vnox_gateway::domain::config::Config = toml::from_str(&text).map_err(|e| anyhow::anyhow!("invalid config: {e}"))?; - info!( - "VNOX Server starting — node: {}", - cfg.node.name - ); + info!("VNOX Server starting — node: {}", cfg.node.name); info!("Gateway TCP: {}", cfg.gateway.bind); info!("Voice UDP: {}", cfg.voice.bind); @@ -59,9 +52,8 @@ async fn main() -> Result<()> { vnox_voice_node::runner::run_bind(&node_name, &voice_bind).await }); - let gate_handle = tokio::spawn(async move { - vnox_gateway::run(gate_cfg, Some(voice_member_tx)).await - }); + let gate_handle = + tokio::spawn(async move { vnox_gateway::run(gate_cfg, Some(voice_member_tx)).await }); tokio::select! { r = voice_handle => { diff --git a/serverd/src/voice_membership.rs b/serverd/src/voice_membership.rs deleted file mode 100644 index 32a376a..0000000 --- a/serverd/src/voice_membership.rs +++ /dev/null @@ -1,22 +0,0 @@ -use tokio::sync::broadcast; - -#[derive(Clone, Debug)] -pub enum VoiceMembershipEvent { - Joined { - channel_id: String, - session_id: String, - user_id: String, - }, - Left { - channel_id: String, - session_id: String, - user_id: String, - }, -} - -pub fn voice_membership_channel() -> ( - broadcast::Sender, - broadcast::Receiver, -) { - broadcast::channel(256) -}