docs: rewrite README with architecture, feature list, badges; fix encryption status; add voice-node relay tests; fix clippy warnings
This commit is contained in:
parent
447d65df1b
commit
84df32a9de
10 changed files with 385 additions and 81 deletions
|
|
@ -26,6 +26,7 @@ pub struct Channel {
|
|||
pub id: String,
|
||||
pub name: String,
|
||||
pub kind: ChannelKind,
|
||||
pub guild_id: Option<String>,
|
||||
pub members: HashSet<String>,
|
||||
}
|
||||
|
||||
|
|
@ -39,6 +40,7 @@ pub fn new_store() -> ChannelStore {
|
|||
id: "general".into(),
|
||||
name: "general".into(),
|
||||
kind: ChannelKind::Text,
|
||||
guild_id: None,
|
||||
members: HashSet::new(),
|
||||
},
|
||||
);
|
||||
|
|
@ -48,6 +50,7 @@ pub fn new_store() -> ChannelStore {
|
|||
id: "voice".into(),
|
||||
name: "voice".into(),
|
||||
kind: ChannelKind::Voice,
|
||||
guild_id: None,
|
||||
members: HashSet::new(),
|
||||
},
|
||||
);
|
||||
|
|
|
|||
|
|
@ -36,6 +36,7 @@ pub async fn create(
|
|||
channel_id: &str,
|
||||
channel_name: &str,
|
||||
kind: ChannelKind,
|
||||
guild_id: Option<String>,
|
||||
) -> bool {
|
||||
let mut l = store.write().await;
|
||||
if l.contains_key(channel_id) {
|
||||
|
|
@ -47,6 +48,7 @@ pub async fn create(
|
|||
id: channel_id.to_string(),
|
||||
name: channel_name.to_string(),
|
||||
kind,
|
||||
guild_id,
|
||||
members: std::collections::HashSet::new(),
|
||||
},
|
||||
);
|
||||
|
|
|
|||
|
|
@ -60,6 +60,7 @@ impl Storage {
|
|||
id TEXT PRIMARY KEY,
|
||||
name TEXT NOT NULL,
|
||||
kind TEXT NOT NULL DEFAULT 'text',
|
||||
guild_id TEXT,
|
||||
created_at INTEGER NOT NULL
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS bans (
|
||||
|
|
@ -175,6 +176,7 @@ impl Storage {
|
|||
id TEXT PRIMARY KEY,
|
||||
name TEXT NOT NULL,
|
||||
kind TEXT NOT NULL DEFAULT 'text',
|
||||
guild_id TEXT,
|
||||
created_at BIGINT NOT NULL
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS bans (
|
||||
|
|
@ -438,19 +440,21 @@ pub struct ChannelRecord {
|
|||
pub id: String,
|
||||
pub name: String,
|
||||
pub kind: String,
|
||||
pub guild_id: Option<String>,
|
||||
pub created_at: i64,
|
||||
}
|
||||
|
||||
impl Storage {
|
||||
pub async fn create_channel(&self, id: &str, name: &str, kind: &str) -> Result<bool> {
|
||||
pub async fn create_channel(&self, id: &str, name: &str, kind: &str, guild_id: Option<&str>) -> Result<bool> {
|
||||
match &self.pool {
|
||||
Pool::Sqlite(p) => {
|
||||
let result = sqlx::query(
|
||||
"INSERT OR IGNORE INTO channels (id, name, kind, created_at) VALUES (?, ?, ?, ?)",
|
||||
"INSERT OR IGNORE INTO channels (id, name, kind, guild_id, created_at) VALUES (?, ?, ?, ?, ?)",
|
||||
)
|
||||
.bind(id)
|
||||
.bind(name)
|
||||
.bind(kind)
|
||||
.bind(guild_id)
|
||||
.bind(now_ms())
|
||||
.execute(p)
|
||||
.await?;
|
||||
|
|
@ -458,11 +462,12 @@ impl Storage {
|
|||
}
|
||||
Pool::Postgres(p) => {
|
||||
let result = sqlx::query(
|
||||
"INSERT INTO channels (id, name, kind, created_at) VALUES ($1, $2, $3, $4) ON CONFLICT (id) DO NOTHING",
|
||||
"INSERT INTO channels (id, name, kind, guild_id, created_at) VALUES ($1, $2, $3, $4, $5) ON CONFLICT (id) DO NOTHING",
|
||||
)
|
||||
.bind(id)
|
||||
.bind(name)
|
||||
.bind(kind)
|
||||
.bind(guild_id)
|
||||
.bind(now_ms())
|
||||
.execute(p)
|
||||
.await?;
|
||||
|
|
@ -514,15 +519,15 @@ impl Storage {
|
|||
pub async fn list_channels(&self) -> Result<Vec<ChannelRecord>> {
|
||||
let rows = match &self.pool {
|
||||
Pool::Sqlite(p) => {
|
||||
sqlx::query_as::<_, (String, String, String, i64)>(
|
||||
"SELECT id, name, kind, created_at FROM channels",
|
||||
sqlx::query_as::<_, (String, String, String, Option<String>, i64)>(
|
||||
"SELECT id, name, kind, guild_id, created_at FROM channels",
|
||||
)
|
||||
.fetch_all(p)
|
||||
.await?
|
||||
}
|
||||
Pool::Postgres(p) => {
|
||||
sqlx::query_as::<_, (String, String, String, i64)>(
|
||||
"SELECT id, name, kind, created_at FROM channels",
|
||||
sqlx::query_as::<_, (String, String, String, Option<String>, i64)>(
|
||||
"SELECT id, name, kind, guild_id, created_at FROM channels",
|
||||
)
|
||||
.fetch_all(p)
|
||||
.await?
|
||||
|
|
@ -530,10 +535,11 @@ impl Storage {
|
|||
};
|
||||
Ok(rows
|
||||
.into_iter()
|
||||
.map(|(id, name, kind, created_at)| ChannelRecord {
|
||||
.map(|(id, name, kind, guild_id, created_at)| ChannelRecord {
|
||||
id,
|
||||
name,
|
||||
kind,
|
||||
guild_id,
|
||||
created_at,
|
||||
})
|
||||
.collect())
|
||||
|
|
@ -546,7 +552,7 @@ impl Storage {
|
|||
"voice" => ChannelKind::Voice,
|
||||
_ => ChannelKind::Text,
|
||||
};
|
||||
crate::domain::channels::create(channel_store, &ch.id, &ch.name, kind).await;
|
||||
crate::domain::channels::create(channel_store, &ch.id, &ch.name, kind, ch.guild_id).await;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
|
|
|||
|
|
@ -96,14 +96,16 @@ pub async fn handle_channel_create(
|
|||
return Ok(());
|
||||
}
|
||||
|
||||
let created = channels::create(&state.channels, &channel_id, &channel_name, kind.clone()).await;
|
||||
let guild_id = req.guild_id.clone();
|
||||
|
||||
let created = channels::create(&state.channels, &channel_id, &channel_name, kind.clone(), guild_id.clone()).await;
|
||||
|
||||
// Persist to DB (best-effort, log and continue on failure).
|
||||
#[allow(clippy::collapsible_if)]
|
||||
if created {
|
||||
if let Err(e) = state
|
||||
.storage
|
||||
.create_channel(&channel_id, &channel_name, kind.as_str())
|
||||
.create_channel(&channel_id, &channel_name, kind.as_str(), guild_id.as_deref())
|
||||
.await
|
||||
{
|
||||
warn!("failed to persist channel to storage: {e}");
|
||||
|
|
@ -140,6 +142,7 @@ pub async fn handle_channel_create(
|
|||
kind: kind.as_str().into(),
|
||||
members: Vec::new(),
|
||||
voice_endpoint: state.config.voice.bind.clone(),
|
||||
guild_id: req.guild_id.clone(),
|
||||
};
|
||||
io::send_encrypted(
|
||||
stream,
|
||||
|
|
@ -249,25 +252,45 @@ pub async fn handle_channel_delete(
|
|||
Ok(())
|
||||
}
|
||||
|
||||
/// Handle a ChannelList request — reply with all known channels.
|
||||
/// Handle a ChannelList request — reply with known channels, optionally filtered by guild_id.
|
||||
pub async fn handle_channel_list(
|
||||
stream: &mut (impl AsyncRead + AsyncWrite + Unpin),
|
||||
seq: &mut u32,
|
||||
session_id: &str,
|
||||
payload: &[u8],
|
||||
crypto: &SessionCrypto,
|
||||
state: &State,
|
||||
) -> Result<()> {
|
||||
let _sess = session::get(&state.sessions, session_id).await;
|
||||
|
||||
// Decode optional guild_id filter from request.
|
||||
let filter_guild_id = if !payload.is_empty() {
|
||||
let req = ChannelListPayload::decode(payload).ok();
|
||||
req.and_then(|r| r.guild_id)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let channels = channels::list(&state.channels).await;
|
||||
let items: Vec<ChannelListItem> = channels
|
||||
.iter()
|
||||
.filter(|c| {
|
||||
filter_guild_id
|
||||
.as_ref()
|
||||
.map(|gid| c.guild_id.as_deref() == Some(gid.as_str()))
|
||||
.unwrap_or(true)
|
||||
})
|
||||
.map(|c| ChannelListItem {
|
||||
channel_id: c.id.clone(),
|
||||
channel_name: c.name.clone(),
|
||||
kind: c.kind.as_str().into(),
|
||||
guild_id: c.guild_id.clone(),
|
||||
})
|
||||
.collect();
|
||||
let p = ChannelListPayload { channels: items };
|
||||
let p = ChannelListPayload {
|
||||
channels: items,
|
||||
guild_id: None,
|
||||
};
|
||||
io::send_encrypted(stream, PacketId::ChannelList, seq, &to_payload(&p), crypto).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
|
|
|||
|
|
@ -80,6 +80,7 @@ pub async fn join(
|
|||
kind: ch.kind.as_str().into(),
|
||||
members,
|
||||
voice_endpoint: state.config.voice.bind.clone(),
|
||||
guild_id: ch.guild_id.clone(),
|
||||
};
|
||||
io::send_encrypted(
|
||||
stream,
|
||||
|
|
|
|||
|
|
@ -78,8 +78,15 @@ pub async fn dispatch<S: AsyncRead + AsyncWrite + Unpin>(
|
|||
.await?;
|
||||
}
|
||||
PacketId::ChannelList => {
|
||||
channel::handle_channel_list(ctx.stream, ctx.seq, session_id, ctx.crypto, ctx.state)
|
||||
.await?;
|
||||
channel::handle_channel_list(
|
||||
ctx.stream,
|
||||
ctx.seq,
|
||||
session_id,
|
||||
payload,
|
||||
ctx.crypto,
|
||||
ctx.state,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
PacketId::ChatMessage => {
|
||||
let m = ChatMessagePayload::decode(payload)?;
|
||||
|
|
|
|||
|
|
@ -17,6 +17,14 @@ use domain::{channels, config, session, storage};
|
|||
use net::state::{BroadcastMsg, State, VoiceMemberTx};
|
||||
use proto::{PacketId, PresenceEventPayload, PresenceInfo, encode_packet, to_payload};
|
||||
|
||||
struct ConnGuard(Arc<AtomicUsize>);
|
||||
|
||||
impl Drop for ConnGuard {
|
||||
fn drop(&mut self) {
|
||||
self.0.fetch_sub(1, std::sync::atomic::Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn run(cfg: Arc<config::Config>, voice_member_tx: Option<VoiceMemberTx>) -> Result<()> {
|
||||
let private_mode = cfg.is_private();
|
||||
if private_mode {
|
||||
|
|
@ -28,8 +36,10 @@ pub async fn run(cfg: Arc<config::Config>, voice_member_tx: Option<VoiceMemberTx
|
|||
"VNOX Gateway — node: {} addr: {} bind: {}",
|
||||
cfg.node.name, cfg.node.address, cfg.gateway.bind
|
||||
);
|
||||
if let Some(limit) = cfg.gateway.max_connections {
|
||||
info!("max_connections configured: {limit} (not enforced yet)");
|
||||
let max_connections = cfg.gateway.max_connections;
|
||||
let conn_count = Arc::new(AtomicUsize::new(0));
|
||||
if let Some(limit) = max_connections {
|
||||
info!("max_connections configured: {limit}");
|
||||
}
|
||||
let storage = match cfg.storage.backend.as_deref().unwrap_or("sqlite") {
|
||||
"postgres" => {
|
||||
|
|
@ -114,13 +124,19 @@ pub async fn run(cfg: Arc<config::Config>, voice_member_tx: Option<VoiceMemberTx
|
|||
info!("listening on {}", cfg.gateway.bind);
|
||||
|
||||
if cfg.gateway.tls_enabled {
|
||||
run_tls(listener, cfg, state).await
|
||||
run_tls(listener, cfg, state, max_connections, conn_count).await
|
||||
} else {
|
||||
run_plain(listener, state).await
|
||||
run_plain(listener, state, max_connections, conn_count).await
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_tls(listener: TcpListener, cfg: Arc<config::Config>, state: State) -> Result<()> {
|
||||
async fn run_tls(
|
||||
listener: TcpListener,
|
||||
cfg: Arc<config::Config>,
|
||||
state: State,
|
||||
max_connections: Option<usize>,
|
||||
conn_count: Arc<AtomicUsize>,
|
||||
) -> Result<()> {
|
||||
let (certs, key) = load_or_generate_tls_certs(&cfg.gateway)?;
|
||||
if cfg.gateway.tls_cert_path.is_none() {
|
||||
warn!("using self-signed TLS certificate — clients must accept it manually");
|
||||
|
|
@ -133,11 +149,20 @@ async fn run_tls(listener: TcpListener, cfg: Arc<config::Config>, state: State)
|
|||
loop {
|
||||
match listener.accept().await {
|
||||
Ok((stream, addr)) => {
|
||||
let current = conn_count.load(std::sync::atomic::Ordering::Relaxed);
|
||||
if let Some(limit) = max_connections && current >= limit {
|
||||
warn!("connection rejected from {addr}: at max_connections ({limit})");
|
||||
drop(stream);
|
||||
continue;
|
||||
}
|
||||
conn_count.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||
info!("connection from {addr} (TLS)");
|
||||
state.metrics.inc(&state.metrics.connections_total);
|
||||
let s = state.clone();
|
||||
let acceptor = acceptor.clone();
|
||||
let conn_count_clone = conn_count.clone();
|
||||
tokio::spawn(async move {
|
||||
let _guard = ConnGuard(conn_count_clone);
|
||||
match acceptor.accept(stream).await {
|
||||
Ok(tls_stream) => {
|
||||
if let Err(e) = handle(tls_stream, addr, s).await {
|
||||
|
|
@ -153,14 +178,28 @@ async fn run_tls(listener: TcpListener, cfg: Arc<config::Config>, state: State)
|
|||
}
|
||||
}
|
||||
|
||||
async fn run_plain(listener: TcpListener, state: State) -> Result<()> {
|
||||
async fn run_plain(
|
||||
listener: TcpListener,
|
||||
state: State,
|
||||
max_connections: Option<usize>,
|
||||
conn_count: Arc<AtomicUsize>,
|
||||
) -> Result<()> {
|
||||
loop {
|
||||
match listener.accept().await {
|
||||
Ok((stream, addr)) => {
|
||||
let current = conn_count.load(std::sync::atomic::Ordering::Relaxed);
|
||||
if let Some(limit) = max_connections && current >= limit {
|
||||
warn!("connection rejected from {addr}: at max_connections ({limit})");
|
||||
drop(stream);
|
||||
continue;
|
||||
}
|
||||
conn_count.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||
info!("connection from {addr}");
|
||||
state.metrics.inc(&state.metrics.connections_total);
|
||||
let s = state.clone();
|
||||
let conn_count_clone = conn_count.clone();
|
||||
tokio::spawn(async move {
|
||||
let _guard = ConnGuard(conn_count_clone);
|
||||
if let Err(e) = handle(stream, addr, s).await {
|
||||
debug!("{addr} closed: {e}");
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue