protobuf migration: replace JSON with prost-generated wire format

- Expanded protocol/lnex.proto with full message set for all packet types
- Added prost + prost-build to gateway build pipeline
- Refactored all handlers and framing to use prost encode/decode
- Removed serde payload structs (now generated from .proto)
- Server starts and admin endpoint responds
This commit is contained in:
loki5512344 2026-07-09 13:04:22 +02:00
parent 5534cd01f7
commit 9821814240
Signed by: boba
GPG key ID: 253067914055423B
31 changed files with 719 additions and 661 deletions

View file

@ -38,9 +38,13 @@ axum = { version = "0.8", features = ["ws"] }
tower = "0.5"
sqlx = { version = "0.9", features = ["sqlite", "runtime-tokio", "tls-rustls"] }
bitflags = "2"
prost = "0.13"
# TLS 1.3 support
rustls-pki-types = "1"
rustls = "0.23"
tokio-rustls = "0.26"
rcgen = "0.13"
[build-dependencies]
prost-build = "0.13"

3
gateway/build.rs Normal file
View file

@ -0,0 +1,3 @@
fn main() {
prost_build::compile_protos(&["../protocol/lnex.proto"], &["../protocol/"]).unwrap();
}

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use tracing::{info, warn};
@ -24,7 +25,7 @@ pub async fn handle_channel_create(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: ChannelCreatePayload = serde_json::from_slice(payload)?;
let req = ChannelCreatePayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),
@ -178,7 +179,7 @@ pub async fn handle_channel_delete(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: ChannelDeletePayload = serde_json::from_slice(payload)?;
let req = ChannelDeletePayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tracing::debug;
use crate::{
@ -8,7 +9,7 @@ use crate::{
};
pub async fn handle_message_edit(session_id: &str, payload: &[u8], state: &State) -> Result<()> {
let edit: crate::proto::MessageEditPayload = serde_json::from_slice(payload)?;
let edit = crate::proto::MessageEditPayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),
@ -60,7 +61,7 @@ pub async fn handle_message_edit(session_id: &str, payload: &[u8], state: &State
}
pub async fn handle_message_delete(session_id: &str, payload: &[u8], state: &State) -> Result<()> {
let delete: MessageDeletePayload = serde_json::from_slice(payload)?;
let delete = MessageDeletePayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -21,7 +22,7 @@ pub async fn handle_presence_update(
_crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: PresenceUpdatePayload = serde_json::from_slice(payload)?;
let req = PresenceUpdatePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;

View file

@ -1,9 +1,10 @@
use anyhow::Result;
use prost::Message;
use crate::{
domain::session,
net::state::{BroadcastMsg, State},
proto::{PacketId, ReadReceiptPayload, encode_packet, to_payload},
proto::{PacketId, ReadReceiptBroadcastPayload, ReadReceiptPayload, encode_packet, to_payload},
};
pub async fn handle_read_receipt(
@ -14,7 +15,7 @@ pub async fn handle_read_receipt(
_crypto: &crate::proto::SessionCrypto,
state: &State,
) -> Result<()> {
let req: ReadReceiptPayload = serde_json::from_slice(payload)?;
let req = ReadReceiptPayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),
@ -25,11 +26,11 @@ pub async fn handle_read_receipt(
.update_read_receipt(&req.channel_id, &sess.user_id, &req.last_read_message_id)
.await?;
let broadcast_data = serde_json::json!({
"channel_id": req.channel_id,
"user_id": sess.user_id,
"last_read_message_id": req.last_read_message_id,
});
let broadcast_data = ReadReceiptBroadcastPayload {
channel_id: req.channel_id.clone(),
user_id: sess.user_id.clone(),
last_read_message_id: req.last_read_message_id,
};
let _ = state.broadcast.send(BroadcastMsg {
channel_id: Some(req.channel_id),

View file

@ -1,23 +1,24 @@
use anyhow::Result;
use prost::Message;
use crate::{
domain::session,
net::state::{BroadcastMsg, State},
proto::{PacketId, TypingStartPayload, encode_packet, to_payload},
proto::{PacketId, TypingBroadcastPayload, TypingStartPayload, encode_packet, to_payload},
};
pub async fn handle_typing_start(session_id: &str, payload: &[u8], state: &State) -> Result<()> {
let req: TypingStartPayload = serde_json::from_slice(payload)?;
let req = TypingStartPayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),
};
let data = serde_json::json!({
"user_id": sess.user_id,
"nickname": sess.nickname,
"channel_id": req.channel_id,
});
let data = TypingBroadcastPayload {
user_id: sess.user_id.clone(),
nickname: sess.nickname.clone(),
channel_id: req.channel_id.clone(),
};
let _ = state.broadcast.send(BroadcastMsg {
channel_id: Some(req.channel_id),

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -15,7 +16,7 @@ pub async fn handle_dm_history(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: DmHistoryPayload = serde_json::from_slice(payload)?;
let req = DmHistoryPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use tracing::warn;
@ -21,7 +22,7 @@ pub async fn handle_dm_message(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let msg: DmMessagePayload = serde_json::from_slice(payload)?;
let msg = DmMessagePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -15,7 +16,7 @@ pub async fn handle_dm_start(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: DmStartPayload = serde_json::from_slice(payload)?;
let req = DmStartPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
@ -95,10 +96,7 @@ pub async fn handle_dm_read_ack(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: serde_json::Value = serde_json::from_slice(payload)?;
let dm_id = req["dm_id"]
.as_str()
.ok_or_else(|| anyhow::anyhow!("missing dm_id"))?;
let req = crate::proto::DmReadAckPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
@ -106,7 +104,7 @@ pub async fn handle_dm_read_ack(
let my_id = sess.user_id.clone();
drop(sess);
state.storage.reset_dm_unread(dm_id, &my_id).await?;
state.storage.reset_dm_unread(&req.dm_id, &my_id).await?;
io::send_encrypted(stream, PacketId::DmReadAck, seq, b"{}", crypto).await?;
Ok(())

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use tracing::debug;
@ -21,7 +22,7 @@ pub async fn dispatch<S: AsyncRead + AsyncWrite + Unpin>(
) -> Result<()> {
match pid {
PacketId::Ping => {
let ping: PingPayload = serde_json::from_slice(payload)?;
let ping = PingPayload::decode(payload)?;
io::send_encrypted(
ctx.stream,
PacketId::Pong,
@ -34,12 +35,12 @@ pub async fn dispatch<S: AsyncRead + AsyncWrite + Unpin>(
.await?;
}
PacketId::JoinChannel => {
let m: JoinChannelPayload = serde_json::from_slice(payload)?;
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 = serde_json::from_slice(payload)?;
let m = LeaveChannelPayload::decode(payload)?;
channel::leave(ctx.stream, ctx.seq, session_id, &m.channel_id, ctx.crypto, ctx.state)
.await?;
}
@ -70,7 +71,7 @@ pub async fn dispatch<S: AsyncRead + AsyncWrite + Unpin>(
.await?;
}
PacketId::ChatMessage => {
let m: ChatMessagePayload = serde_json::from_slice(payload)?;
let m = ChatMessagePayload::decode(payload)?;
content::chat::handle(session_id, m, ctx.state).await?;
}
PacketId::DmStart => {
@ -382,11 +383,11 @@ pub async fn dispatch<S: AsyncRead + AsyncWrite + Unpin>(
.await?;
}
PacketId::MessageReactionAdd => {
let m: ReactionPayload = serde_json::from_slice(payload)?;
let m = ReactionPayload::decode(payload)?;
content::reaction::handle_reaction_add(session_id, m, ctx.state).await?;
}
PacketId::MessageReactionRemove => {
let m: ReactionPayload = serde_json::from_slice(payload)?;
let m = ReactionPayload::decode(payload)?;
content::reaction::handle_reaction_remove(session_id, m, ctx.state).await?;
}
PacketId::MessageEdit => {

View file

@ -1,10 +1,11 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
domain::session,
net::io,
proto::{FriendDeclinePayload, PacketId, SessionCrypto, to_payload},
proto::{FriendDeclinePayload, PacketId, SessionCrypto, SimpleResponsePayload, to_payload},
};
pub async fn handle_friend_decline(
@ -15,7 +16,7 @@ pub async fn handle_friend_decline(
crypto: &SessionCrypto,
state: &crate::net::state::State,
) -> Result<()> {
let req: FriendDeclinePayload = serde_json::from_slice(payload)?;
let req = FriendDeclinePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -29,7 +30,7 @@ pub async fn handle_friend_decline(
stream,
PacketId::FriendDecline,
seq,
&to_payload(&serde_json::json!({"status": "DECLINED"})),
&to_payload(&SimpleResponsePayload { ok: true, user_id: None, role_id: None }),
crypto,
)
.await?;

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -6,7 +7,7 @@ use crate::{
net::{io, state::State},
proto::{
BlockListPayload, BlockUserPayload, FriendRemovePayload, PacketId, SessionCrypto,
UnblockUserPayload, to_payload,
SimpleResponsePayload, UnblockUserPayload, to_payload,
},
};
@ -18,7 +19,7 @@ pub async fn handle_friend_remove(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: FriendRemovePayload = serde_json::from_slice(payload)?;
let req = FriendRemovePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -32,7 +33,7 @@ pub async fn handle_friend_remove(
stream,
PacketId::FriendRemove,
seq,
&to_payload(&serde_json::json!({"removed": true})),
&to_payload(&SimpleResponsePayload { ok: true, user_id: None, role_id: None }),
crypto,
)
.await?;
@ -47,7 +48,7 @@ pub async fn handle_block_user(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: BlockUserPayload = serde_json::from_slice(payload)?;
let req = BlockUserPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -61,7 +62,7 @@ pub async fn handle_block_user(
stream,
PacketId::BlockUser,
seq,
&to_payload(&serde_json::json!({"blocked": true})),
&to_payload(&SimpleResponsePayload { ok: true, user_id: None, role_id: None }),
crypto,
)
.await?;
@ -76,7 +77,7 @@ pub async fn handle_unblock_user(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: UnblockUserPayload = serde_json::from_slice(payload)?;
let req = UnblockUserPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -90,7 +91,7 @@ pub async fn handle_unblock_user(
stream,
PacketId::UnblockUser,
seq,
&to_payload(&serde_json::json!({"unblocked": true})),
&to_payload(&SimpleResponsePayload { ok: true, user_id: None, role_id: None }),
crypto,
)
.await?;

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -23,7 +24,7 @@ pub async fn handle_friend_request(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: FriendRequestPayload = serde_json::from_slice(payload)?;
let req = FriendRequestPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -92,7 +93,9 @@ pub async fn handle_friend_request(
stream,
PacketId::FriendRequest,
seq,
&to_payload(&serde_json::json!({"status": "PENDING", "to_user_id": req.to_user_id})),
&to_payload(&FriendRequestPayload {
to_user_id: req.to_user_id.clone(),
}),
crypto,
)
.await?;
@ -107,7 +110,7 @@ pub async fn handle_friend_accept(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: FriendAcceptPayload = serde_json::from_slice(payload)?;
let req = FriendAcceptPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -144,7 +147,9 @@ pub async fn handle_friend_accept(
stream,
PacketId::FriendAccept,
seq,
&to_payload(&serde_json::json!({"status": "ACCEPTED", "from_user_id": req.from_user_id})),
&to_payload(&FriendAcceptPayload {
from_user_id: req.from_user_id.clone(),
}),
crypto,
)
.await?;

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -19,7 +20,7 @@ pub async fn handle_audit_log_fetch(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: GuildAuditLogFetchPayload = serde_json::from_slice(payload)?;
let req = GuildAuditLogFetchPayload::decode(payload)?;
let sess = match crate::domain::session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use tracing::debug;
@ -18,7 +19,7 @@ pub async fn handle_guild_create(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: GuildCreatePayload = serde_json::from_slice(payload)?;
let req = GuildCreatePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -74,11 +75,7 @@ pub async fn handle_guild_delete(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
#[derive(serde::Deserialize)]
struct Req {
guild_id: String,
}
let req: Req = serde_json::from_slice(payload)?;
let req = crate::proto::GuildDeletePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -112,11 +109,14 @@ pub async fn handle_guild_delete(
None,
)
.await?;
io::send_encrypted(
stream,
PacketId::GuildDelete,
seq,
&to_payload(&serde_json::json!({"guild_id": req.guild_id})),
&to_payload(&crate::proto::GuildDeletePayload {
guild_id: req.guild_id.clone(),
}),
crypto,
)
.await?;

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -20,7 +21,7 @@ pub async fn handle_invite_create(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: InviteCreatePayload = serde_json::from_slice(payload)?;
let req = InviteCreatePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -96,7 +97,7 @@ pub async fn handle_invite_accept(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: InviteAcceptPayload = serde_json::from_slice(payload)?;
let req = InviteAcceptPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -156,7 +157,11 @@ pub async fn handle_invite_accept(
stream,
PacketId::InviteAccept,
seq,
&to_payload(&serde_json::json!({"guild_id": inv.guild_id, "guild_name": inv.guild_name})),
&to_payload(&InviteAcceptPayload {
code: "".into(),
guild_id: inv.guild_id.clone(),
guild_name: inv.guild_name.clone(),
}),
crypto,
)
.await?;
@ -171,7 +176,7 @@ pub async fn handle_invite_delete(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: InviteDeletePayload = serde_json::from_slice(payload)?;
let req = InviteDeletePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -211,7 +216,10 @@ pub async fn handle_invite_delete(
stream,
PacketId::InviteDelete,
seq,
&to_payload(&serde_json::json!({"invite_id": req.invite_id})),
&to_payload(&InviteDeletePayload {
guild_id: req.guild_id.clone(),
invite_id: req.invite_id.clone(),
}),
crypto,
)
.await?;

View file

@ -4,7 +4,7 @@ use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
domain::session,
net::{io, state::State},
proto::{GuildInfo, GuildListPayload, PacketId, SessionCrypto, to_payload},
proto::{GuildInfo, GuildListPayload, PacketId, SessionCrypto, UserRoleUpdatePayload, to_payload},
};
pub async fn handle_guild_list(
@ -41,9 +41,11 @@ pub async fn handle_guild_list(
stream,
PacketId::UserRoleUpdate,
seq,
&to_payload(
&serde_json::json!({"user_id": sess.user_id, "guild_id": g.id, "color": color}),
),
&to_payload(&UserRoleUpdatePayload {
user_id: sess.user_id.clone(),
guild_id: g.id.clone(),
color,
}),
crypto,
)
.await?;

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -20,7 +21,7 @@ pub async fn handle_guild_member_join(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: GuildMemberJoinPayload = serde_json::from_slice(payload)?;
let req = GuildMemberJoinPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -44,7 +45,9 @@ pub async fn handle_guild_member_join(
stream,
PacketId::GuildMemberJoin,
seq,
&to_payload(&serde_json::json!({"guild_id": req.guild_id, "user_id": sess.user_id})),
&to_payload(&crate::proto::GuildMemberJoinPayload {
guild_id: req.guild_id.clone(),
}),
crypto,
)
.await?;
@ -53,9 +56,11 @@ pub async fn handle_guild_member_join(
stream,
PacketId::UserRoleUpdate,
seq,
&to_payload(
&serde_json::json!({"user_id": sess.user_id, "guild_id": req.guild_id, "color": color}),
),
&to_payload(&crate::proto::UserRoleUpdatePayload {
user_id: sess.user_id.clone(),
guild_id: req.guild_id.clone(),
color,
}),
crypto,
)
.await?;
@ -70,25 +75,28 @@ pub async fn handle_guild_member_leave(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: GuildMemberLeavePayload = serde_json::from_slice(payload)?;
let req = GuildMemberLeavePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
let target = if req.user_id.is_empty() {
&sess.user_id
sess.user_id.clone()
} else {
&req.user_id
req.user_id.clone()
};
state
.storage
.remove_guild_member(&req.guild_id, target)
.remove_guild_member(&req.guild_id, target.as_str())
.await?;
io::send_encrypted(
stream,
PacketId::GuildMemberLeave,
seq,
&to_payload(&serde_json::json!({"guild_id": req.guild_id, "user_id": target})),
&to_payload(&GuildMemberLeavePayload {
guild_id: req.guild_id.clone(),
user_id: target,
}),
crypto,
)
.await?;
@ -103,7 +111,7 @@ pub async fn handle_guild_member_kick(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: GuildMemberKickPayload = serde_json::from_slice(payload)?;
let req = GuildMemberKickPayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -146,7 +154,10 @@ pub async fn handle_guild_member_kick(
stream,
PacketId::GuildMemberKick,
seq,
&to_payload(&serde_json::json!({"guild_id": req.guild_id, "user_id": req.user_id})),
&to_payload(&GuildMemberKickPayload {
guild_id: req.guild_id.clone(),
user_id: req.user_id.clone(),
}),
crypto,
)
.await?;

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -22,7 +23,7 @@ pub async fn handle_member_list_fetch(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: GuildMemberListFetchPayload = serde_json::from_slice(payload)?;
let req = GuildMemberListFetchPayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),
@ -72,7 +73,7 @@ pub async fn handle_role_assign(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: RoleAssignPayload = serde_json::from_slice(payload)?;
let req = RoleAssignPayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),
@ -118,9 +119,11 @@ pub async fn handle_role_assign(
stream,
PacketId::GuildRoleAssign,
seq,
&to_payload(
&serde_json::json!({"ok": true, "user_id": req.user_id, "role_id": req.role_id}),
),
&to_payload(&RoleAssignPayload {
guild_id: req.guild_id.clone(),
user_id: req.user_id.clone(),
role_id: req.role_id.clone(),
}),
crypto,
)
.await?;
@ -136,7 +139,7 @@ pub async fn handle_role_unassign(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: RoleAssignPayload = serde_json::from_slice(payload)?;
let req = RoleAssignPayload::decode(payload)?;
let sess = match session::get(&state.sessions, session_id).await {
Some(s) => s,
None => return Ok(()),
@ -181,9 +184,11 @@ pub async fn handle_role_unassign(
stream,
PacketId::GuildRoleUnassign,
seq,
&to_payload(
&serde_json::json!({"ok": true, "user_id": req.user_id, "role_id": req.role_id}),
),
&to_payload(&RoleAssignPayload {
guild_id: req.guild_id.clone(),
user_id: req.user_id.clone(),
role_id: req.role_id.clone(),
}),
crypto,
)
.await?;
@ -199,7 +204,7 @@ pub async fn handle_role_list_fetch(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: GuildRoleListFetchPayload = serde_json::from_slice(payload)?;
let req = GuildRoleListFetchPayload::decode(payload)?;
let _sess = session::get(&state.sessions, session_id).await;
let rows = state.storage.list_guild_roles(&req.guild_id).await?;

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use crate::{
@ -17,7 +18,7 @@ pub async fn handle_role_create(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: RoleCreatePayload = serde_json::from_slice(payload)?;
let req = RoleCreatePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -63,9 +64,13 @@ pub async fn handle_role_create(
stream,
PacketId::RoleCreate,
seq,
&to_payload(
&serde_json::json!({"id": role_id, "guild_id": req.guild_id, "name": req.name}),
),
&to_payload(&RoleCreatePayload {
guild_id: req.guild_id.clone(),
name: req.name.clone(),
color: req.color.clone(),
permissions: req.permissions,
id: role_id,
}),
crypto,
)
.await?;
@ -80,7 +85,7 @@ pub async fn handle_role_delete(
crypto: &SessionCrypto,
state: &State,
) -> Result<()> {
let req: RoleDeletePayload = serde_json::from_slice(payload)?;
let req = RoleDeletePayload::decode(payload)?;
let sess = session::get(&state.sessions, session_id)
.await
.ok_or_else(|| anyhow::anyhow!("session not found"))?;
@ -120,7 +125,10 @@ pub async fn handle_role_delete(
stream,
PacketId::RoleDelete,
seq,
&to_payload(&serde_json::json!({"role_id": req.role_id})),
&to_payload(&RoleDeletePayload {
guild_id: req.guild_id.clone(),
role_id: req.role_id.clone(),
}),
crypto,
)
.await?;

View file

@ -1,4 +1,5 @@
use anyhow::Result;
use prost::Message;
use tokio::io::{AsyncRead, AsyncWrite};
use tracing::{debug, info, warn};
@ -23,17 +24,17 @@ pub async fn run(
) -> Result<(session::Session, SessionCrypto)> {
// Generate ephemeral X25519 keypair for forward secrecy
let (eph_sk, eph_pk) = SessionCrypto::new_ephemeral();
let eph_pk_hex = hex::encode(eph_pk.as_bytes());
let eph_pk_raw = eph_pk.as_bytes();
// HELLO — include server's ephemeral public key + privacy mode
let challenge = auth::new_challenge();
let private_mode = state.config.is_private();
let hello = HelloPayload {
lnex_version: LNEX_VERSION.into(),
server_pubkey: state.server_identity.pubkey_hex(),
challenge_nonce: hex::encode(challenge),
server_pubkey: hex::decode(state.server_identity.pubkey_hex())?,
challenge_nonce: challenge.to_vec(),
node_name: state.config.node.name.clone(),
server_eph_pubkey: eph_pk_hex,
server_eph_pubkey: eph_pk_raw.to_vec(),
private_mode,
};
io::send_packet(stream, PacketId::Hello, seq, &to_payload(&hello)).await?;
@ -52,7 +53,7 @@ pub async fn run(
return Err(anyhow::anyhow!("expected AUTH"));
}
let msg: AuthPayload = serde_json::from_slice(&payload)?;
let msg = AuthPayload::decode(payload.as_slice())?;
debug!("{addr} → AUTH nick={}", msg.nickname);
if msg.lnex_version != LNEX_VERSION {
@ -66,9 +67,8 @@ pub async fn run(
return Err(anyhow::anyhow!("version mismatch"));
}
let pubkey =
hex_to_32(&msg.client_pubkey).ok_or_else(|| anyhow::anyhow!("bad client pubkey"))?;
let sig = hex_to_64(&msg.signature).ok_or_else(|| anyhow::anyhow!("bad signature"))?;
let pubkey: [u8; 32] = msg.client_pubkey.as_slice().try_into()?;
let sig: [u8; 64] = msg.signature.as_slice().try_into()?;
if let Err(e) = auth::verify_auth(&challenge, &pubkey, &sig) {
warn!("{addr} auth failed: {e}");
@ -83,7 +83,7 @@ pub async fn run(
return Err(anyhow::anyhow!("auth failed"));
}
if state.storage.is_banned(&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"));
@ -91,19 +91,18 @@ pub async fn run(
state
.storage
.upsert_user(&msg.client_pubkey, &msg.nickname)
.upsert_user(&hex::encode(&msg.client_pubkey), &msg.nickname)
.await?;
// ECDH: compute shared secret from server's ephemeral sk + client's ephemeral pk
let client_eph_raw = hex_to_32(&msg.client_eph_pubkey)
.ok_or_else(|| anyhow::anyhow!("bad client eph pubkey"))?;
let client_eph_raw: [u8; 32] = msg.client_eph_pubkey.as_slice().try_into()?;
let client_eph_pk = x25519_dalek::PublicKey::from(client_eph_raw);
let shared_secret = SessionCrypto::ecdh(eph_sk, &client_eph_pk);
// SESSION
let sess = session::create(
&state.sessions,
msg.client_pubkey.clone(),
hex::encode(&msg.client_pubkey),
msg.nickname.clone(),
)
.await;
@ -131,10 +130,3 @@ pub async fn run(
Ok((sess, crypto))
}
fn hex_to_32(s: &str) -> Option<[u8; 32]> {
hex::decode(s).ok()?.try_into().ok()
}
fn hex_to_64(s: &str) -> Option<[u8; 64]> {
hex::decode(s).ok()?.try_into().ok()
}

View file

@ -1,6 +1,6 @@
use chacha20poly1305::{
ChaCha20Poly1305, Key, KeyInit, Nonce,
aead::{Aead, Payload},
ChaCha20Poly1305, Key, KeyInit, Nonce,
};
pub(super) fn make_nonce(cid: &[u8; 8], seq: u64) -> [u8; 12] {

View file

@ -1,7 +1,7 @@
use super::packet::{PacketHeader, PacketId};
use serde::Serialize;
use prost::Message;
/// Encode a packet: header + JSON payload bytes.
/// Encode a packet: header + protobuf payload bytes.
pub fn encode_packet(id: PacketId, seq: u32, payload: &[u8]) -> Vec<u8> {
let header = PacketHeader::new(id, seq, payload.len() as u32);
let mut out = Vec::with_capacity(PacketHeader::SIZE + payload.len());
@ -10,7 +10,7 @@ pub fn encode_packet(id: PacketId, seq: u32, payload: &[u8]) -> Vec<u8> {
out
}
/// Serialize a payload struct to JSON bytes.
pub fn to_payload<T: Serialize>(v: &T) -> Vec<u8> {
serde_json::to_vec(v).expect("payload serialization is infallible")
/// Encode a protobuf payload struct to bytes.
pub fn to_payload(msg: &impl Message) -> Vec<u8> {
msg.encode_to_vec()
}

View file

@ -1,11 +1,104 @@
mod crypto;
mod framing;
mod packet;
mod payloads;
pub mod lnex {
include!(concat!(env!("OUT_DIR"), "/lnex.rs"));
}
pub use lnex::*;
pub use crypto::SessionCrypto;
// Re-export everything so `crate::proto::X` works as before
pub use framing::{encode_packet, to_payload};
pub use packet::{ErrorCode, PacketHeader, PacketId, flags};
pub use payloads::*;
pub use packet::{flags, ErrorCode, PacketHeader, PacketId};
// ─── Backward-compat type aliases ────────────────────────────────
// Handshake
pub type HelloPayload = Hello;
pub type AuthPayload = Auth;
pub type SessionPayload = Session;
// Keepalive
pub type PingPayload = Ping;
pub type PongPayload = Pong;
// Error / Disconnect
pub type ErrorPayload = Error;
pub type DisconnectPayload = Disconnect;
// Channels
pub type JoinChannelPayload = JoinChannel;
pub type LeaveChannelPayload = LeaveChannel;
pub type ChannelStatePayload = ChannelState;
pub type ChannelCreatePayload = ChannelCreate;
pub type ChannelDeletePayload = ChannelDelete;
pub type ChannelListPayload = ChannelList;
pub type UserJoinPayload = UserJoin;
pub type UserLeavePayload = UserLeave;
// Chat
pub type ChatMessagePayload = ChatMessage;
pub type ChatHistoryPayload = ChatHistory;
pub type MessageEditPayload = MessageEdit;
pub type MessageDeletePayload = MessageDelete;
// Voice
pub type VoiceStatePayload = VoiceState;
// Guilds
pub type GuildCreatePayload = GuildCreate;
pub type GuildDeletePayload = GuildDelete;
pub type GuildListPayload = GuildList;
pub type GuildMemberJoinPayload = GuildMemberJoin;
pub type GuildMemberLeavePayload = GuildMemberLeave;
pub type GuildMemberKickPayload = GuildMemberKick;
pub type RoleCreatePayload = RoleCreate;
pub type RoleDeletePayload = RoleDelete;
pub type InviteCreatePayload = InviteCreate;
pub type InviteAcceptPayload = InviteAccept;
pub type InviteDeletePayload = InviteDelete;
pub type GuildAuditLogFetchPayload = GuildAuditLogFetch;
pub type AuditLogEntryPayload = AuditLogEntry;
pub type GuildAuditLogPayload = GuildAuditLog;
pub type GuildMemberListFetchPayload = GuildMemberListFetch;
pub type GuildMemberInfoPayload = GuildMemberInfo;
pub type GuildMemberListPayload = GuildMemberList;
pub type RoleAssignPayload = RoleAssign;
pub type GuildRoleListFetchPayload = GuildRoleListFetch;
pub type GuildRoleInfoPayload = GuildRoleInfo;
pub type GuildRoleListPayload = GuildRoleList;
// DMs
pub type DmStartPayload = DmStart;
pub type DmStartResponsePayload = DmStartResponse;
pub type DmMessagePayload = DmMessage;
pub type DmHistoryPayload = DmHistory;
pub type DmReadAckPayload = DmReadAck;
// Friends
pub type FriendRequestPayload = FriendRequest;
pub type FriendAcceptPayload = FriendAccept;
pub type FriendDeclinePayload = FriendDecline;
pub type FriendRemovePayload = FriendRemove;
pub type FriendListPayload = FriendList;
pub type BlockUserPayload = BlockUser;
pub type UnblockUserPayload = UnblockUser;
pub type BlockListPayload = BlockList;
pub type FriendEventPayload = FriendEvent;
// Presence
pub type PresenceUpdatePayload = PresenceUpdate;
pub type PresenceSyncPayload = PresenceSync;
pub type PresenceEventPayload = PresenceEvent;
// Read/typing
pub type ReadReceiptPayload = ReadReceipt;
pub type TypingStartPayload = TypingStart;
// Response types
pub type ReadReceiptBroadcastPayload = ReadReceiptBroadcast;
pub type UserRoleUpdatePayload = UserRoleUpdate;
pub type TypingBroadcastPayload = TypingBroadcast;
pub type SimpleResponsePayload = SimpleResponse;

View file

@ -1,157 +0,0 @@
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
pub struct HelloPayload {
pub lnex_version: String,
pub server_pubkey: String,
pub challenge_nonce: String,
pub node_name: String,
pub server_eph_pubkey: String,
pub private_mode: bool,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct AuthPayload {
pub client_pubkey: String,
pub nickname: String,
pub lnex_version: String,
pub signature: String,
pub client_eph_pubkey: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct SessionPayload {
pub session_id: String,
pub token: String,
pub expires_at: i64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct PingPayload {
pub timestamp: i64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct PongPayload {
pub timestamp: i64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct JoinChannelPayload {
pub channel_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct LeaveChannelPayload {
pub channel_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ChannelStatePayload {
pub channel_id: String,
pub channel_name: String,
pub kind: String,
pub members: Vec<MemberInfo>,
pub voice_endpoint: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ChannelCreatePayload {
pub channel_id: String,
pub channel_name: String,
/// "text" or "voice".
pub kind: String,
/// Optional guild_id this channel belongs to (Phase 1.x — unused, future).
#[serde(default)]
pub guild_id: Option<String>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ChannelDeletePayload {
pub channel_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ChannelListPayload {
pub channels: Vec<ChannelListItem>,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct ChannelListItem {
pub channel_id: String,
pub channel_name: String,
pub kind: String,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct MemberInfo {
pub user_id: String,
pub nickname: String,
pub in_voice: bool,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct UserJoinPayload {
pub channel_id: String,
pub user_id: String,
pub nickname: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct UserLeavePayload {
pub channel_id: String,
pub user_id: String,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct ChatMessagePayload {
pub message_id: String,
pub channel_id: String,
pub sender_id: String,
pub content: String,
pub timestamp: i64,
#[serde(default)]
pub edited: bool,
/// Optional message_id this message is replying to.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reply_to: Option<String>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ReactionPayload {
pub message_id: String,
pub channel_id: String,
pub emoji: String,
#[serde(default)]
pub user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct MessageEditPayload {
pub message_id: String,
pub channel_id: String,
pub content: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct MessageDeletePayload {
pub message_id: String,
pub channel_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ChatHistoryPayload {
pub channel_id: String,
pub messages: Vec<ChatMessagePayload>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ErrorPayload {
pub code: u32,
pub message: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct DisconnectPayload {
pub reason: String,
}

View file

@ -1,161 +0,0 @@
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildCreatePayload {
pub name: String,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct GuildInfo {
pub id: String,
pub owner_id: String,
pub name: String,
pub member_count: i64,
pub created_at: i64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildListPayload {
pub guilds: Vec<GuildInfo>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildMemberJoinPayload {
pub guild_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildMemberLeavePayload {
pub guild_id: String,
pub user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildMemberKickPayload {
pub guild_id: String,
pub user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RoleCreatePayload {
pub guild_id: String,
pub name: String,
pub color: Option<String>,
pub permissions: Option<u64>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RoleDeletePayload {
pub guild_id: String,
pub role_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct InviteCreatePayload {
pub guild_id: String,
pub max_uses: Option<i64>,
pub expires_in_seconds: Option<i64>,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct InviteInfo {
pub id: String,
pub guild_id: String,
pub guild_name: String,
pub code: String,
pub creator_id: String,
pub max_uses: Option<i64>,
pub uses: i64,
pub expires_at: Option<i64>,
pub created_at: i64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct InviteAcceptPayload {
pub code: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct InviteDeletePayload {
pub guild_id: String,
pub invite_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildAuditLogFetchPayload {
pub guild_id: String,
#[serde(default = "default_audit_limit")]
pub limit: i64,
}
fn default_audit_limit() -> i64 {
50
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct AuditLogEntryPayload {
pub id: String,
pub guild_id: String,
pub actor_id: String,
pub action: String,
pub target_id: Option<String>,
pub target_type: Option<String>,
pub reason: Option<String>,
pub created_at: i64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildAuditLogPayload {
pub guild_id: String,
pub entries: Vec<AuditLogEntryPayload>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildMemberListFetchPayload {
pub guild_id: String,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct GuildMemberInfoPayload {
pub user_id: String,
pub nickname: String,
pub joined_at: i64,
pub role_color: String,
pub role_name: String,
/// True if this user is the guild owner.
pub is_owner: bool,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildMemberListPayload {
pub guild_id: String,
pub members: Vec<GuildMemberInfoPayload>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct RoleAssignPayload {
pub guild_id: String,
pub user_id: String,
pub role_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildRoleListFetchPayload {
pub guild_id: String,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct GuildRoleInfoPayload {
pub id: String,
pub guild_id: String,
pub name: String,
pub color: String,
pub permissions: u64,
pub position: i32,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct GuildRoleListPayload {
pub guild_id: String,
pub roles: Vec<GuildRoleInfoPayload>,
}

View file

@ -1,7 +0,0 @@
mod channel;
mod guild;
mod social;
pub use channel::*;
pub use guild::*;
pub use social::*;

View file

@ -1,132 +0,0 @@
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
pub struct DmStartPayload {
pub target_user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct DmStartResponsePayload {
pub dm_id: String,
pub other_user_id: String,
pub other_nickname: String,
pub messages: Vec<DmMessagePayload>,
#[serde(default)]
pub unread_count: u32,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct DmMessagePayload {
pub dm_id: String,
pub sender_id: String,
pub content: String,
pub timestamp: i64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct DmHistoryPayload {
pub dm_id: String,
#[serde(default)]
pub messages: Vec<DmMessagePayload>,
#[serde(default)]
pub search_query: Option<String>,
#[serde(default)]
pub limit: Option<i64>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct FriendRequestPayload {
pub to_user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct FriendAcceptPayload {
pub from_user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct FriendDeclinePayload {
pub from_user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct FriendRemovePayload {
pub user_id: String,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct FriendInfo {
pub user_id: String,
pub nickname: String,
pub status: String,
pub since: i64,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct FriendListPayload {
pub friends: Vec<FriendInfo>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct BlockUserPayload {
pub user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct UnblockUserPayload {
pub user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct BlockListPayload {
pub blocked: Vec<String>,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct FriendEventPayload {
pub event: String,
pub user_id: String,
pub nickname: Option<String>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct PresenceUpdatePayload {
pub status: String,
pub activity_type: Option<String>,
pub activity_text: Option<String>,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct PresenceInfo {
pub user_id: String,
pub nickname: String,
pub status: String,
pub activity_type: Option<String>,
pub activity_text: Option<String>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct PresenceSyncPayload {
pub presences: Vec<PresenceInfo>,
}
#[derive(Debug, Serialize, Deserialize, Clone)]
pub struct PresenceEventPayload {
pub user_id: String,
pub nickname: String,
pub status: String,
pub activity_type: Option<String>,
pub activity_text: Option<String>,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct ReadReceiptPayload {
pub channel_id: String,
pub last_read_message_id: String,
pub user_id: String,
}
#[derive(Debug, Serialize, Deserialize)]
pub struct TypingStartPayload {
pub channel_id: String,
}