Initial commit: split VNOX monorepo into VNOX-Server
Unified server binary (vnox-serverd) combining: - Gateway (TCP, auth, channels, sessions, SQLite) - Voice-node (UDP relay, Opus, jitter buffer) Standalone gateway (vnox-gateway) and voice-node (vnox-voice-node) preserved as sub-crates.
This commit is contained in:
commit
e751e0bf5b
94 changed files with 9771 additions and 0 deletions
26
voice-node/Cargo.toml
Normal file
26
voice-node/Cargo.toml
Normal file
|
|
@ -0,0 +1,26 @@
|
|||
[package]
|
||||
name = "vnox-voice-node"
|
||||
description = "VNOX voice node — UDP relay, Opus, jitter buffer"
|
||||
version.workspace = true
|
||||
edition.workspace = true
|
||||
license.workspace = true
|
||||
authors.workspace = true
|
||||
repository.workspace = true
|
||||
|
||||
[lib]
|
||||
name = "vnox_voice_node"
|
||||
path = "src/lib.rs"
|
||||
|
||||
[[bin]]
|
||||
name = "vnox-voice-node"
|
||||
path = "src/main.rs"
|
||||
|
||||
[dependencies]
|
||||
tokio.workspace = true
|
||||
opus.workspace = true
|
||||
tracing.workspace = true
|
||||
tracing-subscriber.workspace = true
|
||||
anyhow.workspace = true
|
||||
thiserror.workspace = true
|
||||
toml.workspace = true
|
||||
serde.workspace = true
|
||||
75
voice-node/src/jitter/adaptive.rs
Normal file
75
voice-node/src/jitter/adaptive.rs
Normal file
|
|
@ -0,0 +1,75 @@
|
|||
use tracing::info;
|
||||
|
||||
use super::JitterBuffer;
|
||||
|
||||
pub const MIN_TARGET_MS: u32 = 20;
|
||||
pub const MAX_TARGET_MS: u32 = 150;
|
||||
pub const JITTER_WINDOW: usize = 64;
|
||||
const LOSS_INCREASE_THRESH: u32 = 5;
|
||||
const LOSS_DECREASE_THRESH: u32 = 1;
|
||||
const STABLE_WINDOWS_BEFORE_DECREASE: u32 = 10;
|
||||
pub const LOSS_WINDOW_PACKETS: u32 = 50;
|
||||
const ADJUST_STEP_MS: u32 = 5;
|
||||
|
||||
impl JitterBuffer {
|
||||
pub(super) fn observed_jitter_ms(&self) -> u32 {
|
||||
if self.arrival_count < 2 {
|
||||
return 0;
|
||||
}
|
||||
let mut deltas = Vec::with_capacity(self.arrival_count - 1);
|
||||
for i in 0..self.arrival_count - 1 {
|
||||
let idx0 = (self.arrival_idx + JITTER_WINDOW - self.arrival_count + i) % JITTER_WINDOW;
|
||||
let idx1 =
|
||||
(self.arrival_idx + JITTER_WINDOW - self.arrival_count + i + 1) % JITTER_WINDOW;
|
||||
let delta = self.arrival_times[idx1].saturating_sub(self.arrival_times[idx0]);
|
||||
deltas.push(delta);
|
||||
}
|
||||
if deltas.is_empty() {
|
||||
return 0;
|
||||
}
|
||||
deltas.sort_unstable();
|
||||
let p90_idx = ((deltas.len() as f64) * 0.90).ceil() as usize - 1;
|
||||
let p90_idx = p90_idx.min(deltas.len() - 1);
|
||||
deltas[p90_idx] as u32
|
||||
}
|
||||
|
||||
pub(super) fn adapt_target(&mut self) {
|
||||
let loss_pct = if self.total_expected > 0 {
|
||||
(self.loss_count as f64 / self.total_expected as f64 * 100.0) as u32
|
||||
} else {
|
||||
0
|
||||
};
|
||||
let jitter_ms = self.observed_jitter_ms();
|
||||
let jitter_based_target = (jitter_ms * 2).clamp(MIN_TARGET_MS, MAX_TARGET_MS);
|
||||
|
||||
if loss_pct >= LOSS_INCREASE_THRESH {
|
||||
let new_target = (self.target_ms + ADJUST_STEP_MS).min(MAX_TARGET_MS);
|
||||
if new_target != self.target_ms {
|
||||
info!(
|
||||
target_ms = new_target,
|
||||
loss_pct, jitter_ms, "adaptive: increasing target due to packet loss"
|
||||
);
|
||||
self.target_ms = new_target;
|
||||
}
|
||||
self.stable_windows = 0;
|
||||
} else if loss_pct <= LOSS_DECREASE_THRESH {
|
||||
self.stable_windows += 1;
|
||||
if self.stable_windows >= STABLE_WINDOWS_BEFORE_DECREASE {
|
||||
let candidate = self.target_ms.saturating_sub(ADJUST_STEP_MS);
|
||||
let new_target = candidate.max(jitter_based_target).max(MIN_TARGET_MS);
|
||||
if new_target != self.target_ms {
|
||||
info!(
|
||||
target_ms = new_target,
|
||||
loss_pct, jitter_ms, "adaptive: decreasing target, link stable"
|
||||
);
|
||||
self.target_ms = new_target;
|
||||
}
|
||||
self.stable_windows = 0;
|
||||
}
|
||||
} else {
|
||||
self.stable_windows = 0;
|
||||
}
|
||||
self.loss_count = 0;
|
||||
self.total_expected = 0;
|
||||
}
|
||||
}
|
||||
100
voice-node/src/jitter/mod.rs
Normal file
100
voice-node/src/jitter/mod.rs
Normal file
|
|
@ -0,0 +1,100 @@
|
|||
use std::collections::BTreeMap;
|
||||
|
||||
pub mod adaptive;
|
||||
pub mod relay;
|
||||
|
||||
pub use adaptive::{JITTER_WINDOW, LOSS_WINDOW_PACKETS, MAX_TARGET_MS, MIN_TARGET_MS};
|
||||
|
||||
pub struct JitterBuffer {
|
||||
target_ms: u32,
|
||||
adaptive: bool,
|
||||
packets: BTreeMap<u32, BufferedPacket>,
|
||||
last_played: Option<u32>,
|
||||
arrival_times: Vec<u64>,
|
||||
arrival_idx: usize,
|
||||
arrival_count: usize,
|
||||
loss_count: u32,
|
||||
total_expected: u32,
|
||||
stable_windows: u32,
|
||||
}
|
||||
|
||||
impl JitterBuffer {
|
||||
pub fn new(target_ms: u32, adaptive: bool) -> Self {
|
||||
Self {
|
||||
target_ms,
|
||||
adaptive,
|
||||
packets: BTreeMap::new(),
|
||||
last_played: None,
|
||||
arrival_times: vec![0u64; JITTER_WINDOW],
|
||||
arrival_idx: 0,
|
||||
arrival_count: 0,
|
||||
loss_count: 0,
|
||||
total_expected: 0,
|
||||
stable_windows: 0,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn push(&mut self, pkt: BufferedPacket) {
|
||||
let arrived = pkt.arrived_at;
|
||||
self.arrival_times[self.arrival_idx] = arrived;
|
||||
self.arrival_idx = (self.arrival_idx + 1) % JITTER_WINDOW;
|
||||
if self.arrival_count < JITTER_WINDOW {
|
||||
self.arrival_count += 1;
|
||||
}
|
||||
|
||||
let gap = self
|
||||
.packets
|
||||
.keys()
|
||||
.next_back()
|
||||
.map(|&highest| {
|
||||
if pkt.voice_seq > highest {
|
||||
pkt.voice_seq.wrapping_sub(highest).saturating_sub(1)
|
||||
} else {
|
||||
0
|
||||
}
|
||||
})
|
||||
.unwrap_or(0);
|
||||
|
||||
if gap > 0 {
|
||||
self.loss_count += gap;
|
||||
}
|
||||
self.total_expected += 1 + gap;
|
||||
self.packets.insert(pkt.voice_seq, pkt);
|
||||
|
||||
if self.adaptive && self.total_expected >= LOSS_WINDOW_PACKETS {
|
||||
self.adapt_target();
|
||||
}
|
||||
}
|
||||
|
||||
pub fn len(&self) -> usize {
|
||||
self.packets.len()
|
||||
}
|
||||
|
||||
pub fn is_empty(&self) -> bool {
|
||||
self.packets.is_empty()
|
||||
}
|
||||
|
||||
pub fn set_target_ms(&mut self, target_ms: u32) {
|
||||
self.target_ms = target_ms.clamp(MIN_TARGET_MS, MAX_TARGET_MS);
|
||||
}
|
||||
|
||||
pub fn target_ms(&self) -> u32 {
|
||||
self.target_ms
|
||||
}
|
||||
|
||||
pub fn is_adaptive(&self) -> bool {
|
||||
self.adaptive
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct BufferedPacket {
|
||||
pub voice_seq: u32,
|
||||
pub timestamp: u32,
|
||||
pub channel_id: u64,
|
||||
pub opus_data: Vec<u8>,
|
||||
pub arrived_at: u64,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
44
voice-node/src/jitter/relay.rs
Normal file
44
voice-node/src/jitter/relay.rs
Normal file
|
|
@ -0,0 +1,44 @@
|
|||
use super::{BufferedPacket, JitterBuffer};
|
||||
|
||||
impl JitterBuffer {
|
||||
pub fn pop_ready(&mut self, now_ms: u64) -> Option<BufferedPacket> {
|
||||
if self.adaptive {
|
||||
return self.pop_ready_adaptive(now_ms);
|
||||
}
|
||||
self.pop_ready_fixed(now_ms)
|
||||
}
|
||||
|
||||
pub fn has_gap(&self) -> bool {
|
||||
match self.last_played {
|
||||
None => false,
|
||||
Some(last) => {
|
||||
let expected = last.wrapping_add(1);
|
||||
!self.packets.contains_key(&expected) && !self.packets.is_empty()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn pop_ready_fixed(&mut self, now_ms: u64) -> Option<BufferedPacket> {
|
||||
let seq = *self.packets.keys().next()?;
|
||||
let pkt = self.packets.get(&seq)?;
|
||||
let buffered_for = now_ms.saturating_sub(pkt.arrived_at);
|
||||
if buffered_for >= self.target_ms as u64 {
|
||||
let pkt = self.packets.remove(&seq)?;
|
||||
self.last_played = Some(seq);
|
||||
return Some(pkt);
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
fn pop_ready_adaptive(&mut self, now_ms: u64) -> Option<BufferedPacket> {
|
||||
let seq = *self.packets.keys().next()?;
|
||||
let pkt = self.packets.get(&seq)?;
|
||||
let buffered_for = now_ms.saturating_sub(pkt.arrived_at);
|
||||
if buffered_for >= self.target_ms as u64 {
|
||||
let pkt = self.packets.remove(&seq)?;
|
||||
self.last_played = Some(seq);
|
||||
return Some(pkt);
|
||||
}
|
||||
None
|
||||
}
|
||||
}
|
||||
103
voice-node/src/jitter/tests.rs
Normal file
103
voice-node/src/jitter/tests.rs
Normal file
|
|
@ -0,0 +1,103 @@
|
|||
use super::*;
|
||||
|
||||
fn pkt(seq: u32, arrived_at: u64) -> BufferedPacket {
|
||||
BufferedPacket {
|
||||
voice_seq: seq,
|
||||
timestamp: 0,
|
||||
channel_id: 1,
|
||||
opus_data: vec![seq as u8],
|
||||
arrived_at,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pop_ready_waits_for_target_ms() {
|
||||
let mut jb = JitterBuffer::new(20, false);
|
||||
jb.push(pkt(1, 100));
|
||||
assert!(jb.pop_ready(110).is_none());
|
||||
assert_eq!(jb.pop_ready(120).unwrap().voice_seq, 1);
|
||||
assert!(jb.pop_ready(120).is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pop_ready_releases_in_sequence_order() {
|
||||
let mut jb = JitterBuffer::new(0, false);
|
||||
jb.push(pkt(2, 0));
|
||||
jb.push(pkt(1, 0));
|
||||
assert_eq!(jb.pop_ready(0).unwrap().voice_seq, 1);
|
||||
assert_eq!(jb.pop_ready(0).unwrap().voice_seq, 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn has_gap_after_played_seq() {
|
||||
let mut jb = JitterBuffer::new(0, false);
|
||||
jb.push(pkt(1, 0));
|
||||
jb.pop_ready(0);
|
||||
jb.push(pkt(3, 0));
|
||||
assert!(jb.has_gap());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn no_gap_when_next_seq_present() {
|
||||
let mut jb = JitterBuffer::new(0, false);
|
||||
jb.push(pkt(1, 0));
|
||||
jb.pop_ready(0);
|
||||
jb.push(pkt(2, 0));
|
||||
assert!(!jb.has_gap());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn adaptive_mode_increases_on_loss() {
|
||||
let mut jb = JitterBuffer::new(40, true);
|
||||
jb.push(pkt(1, 0));
|
||||
jb.push(pkt(3, 1));
|
||||
for i in 4..=50 {
|
||||
jb.push(pkt(i, i as u64));
|
||||
}
|
||||
assert!(jb.target_ms <= MAX_TARGET_MS);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn set_target_ms_clamps() {
|
||||
let mut jb = JitterBuffer::new(40, false);
|
||||
jb.set_target_ms(0);
|
||||
assert_eq!(jb.target_ms, MIN_TARGET_MS);
|
||||
jb.set_target_ms(500);
|
||||
assert_eq!(jb.target_ms, MAX_TARGET_MS);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn adaptive_mode_respects_flag() {
|
||||
let mut jb = JitterBuffer::new(40, false);
|
||||
jb.push(pkt(1, 0));
|
||||
for i in 2..=100 {
|
||||
jb.push(pkt(i, i as u64));
|
||||
}
|
||||
assert_eq!(jb.target_ms, 40);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn observed_jitter_smooth_timeline() {
|
||||
let mut jb = JitterBuffer::new(40, true);
|
||||
for i in 0..JITTER_WINDOW {
|
||||
jb.push(pkt(i as u32, (i as u64) * 10));
|
||||
}
|
||||
let jitter = jb.observed_jitter_ms();
|
||||
assert!(
|
||||
(8..=12).contains(&jitter),
|
||||
"jitter = {jitter} (expected ~10)"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn observed_jitter_spike() {
|
||||
let mut jb = JitterBuffer::new(40, true);
|
||||
let times = [
|
||||
0, 10, 20, 70, 80, 90, 140, 150, 160, 210, 220, 230, 280, 290, 300,
|
||||
];
|
||||
for (i, &t) in times.iter().enumerate() {
|
||||
jb.push(pkt(i as u32, t));
|
||||
}
|
||||
let jitter = jb.observed_jitter_ms();
|
||||
assert!(jitter >= 40, "jitter = {jitter} (expected >= 40)");
|
||||
}
|
||||
67
voice-node/src/lib.rs
Normal file
67
voice-node/src/lib.rs
Normal file
|
|
@ -0,0 +1,67 @@
|
|||
pub mod jitter;
|
||||
pub mod relay;
|
||||
pub mod runner;
|
||||
use anyhow::Result;
|
||||
use serde::Deserialize;
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct Config {
|
||||
pub node: NodeConfig,
|
||||
pub voice: VoiceConfig,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct NodeConfig {
|
||||
pub name: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub struct VoiceConfig {
|
||||
pub bind: String,
|
||||
}
|
||||
|
||||
pub use runner::run_bind;
|
||||
|
||||
pub fn load_config() -> Result<Config> {
|
||||
let path = std::env::args()
|
||||
.skip_while(|a| a != "--config")
|
||||
.nth(1)
|
||||
.unwrap_or_else(|| "/etc/vnox/config.toml".into());
|
||||
|
||||
let text = std::fs::read_to_string(&path)
|
||||
.map_err(|e| anyhow::anyhow!("cannot read config {path}: {e}"))?;
|
||||
|
||||
toml::from_str(&text).map_err(|e| anyhow::anyhow!("invalid config: {e}"))
|
||||
}
|
||||
|
||||
// ─── Voice packet header ──────────────────────────────────────────────────────
|
||||
|
||||
pub const VOICE_HDR_SIZE: usize = 20;
|
||||
pub const VOICE_PACKET_ID: u16 = 0x0010;
|
||||
pub const MAX_UDP_PACKET: usize = 1472;
|
||||
pub const PLAYOUT_INTERVAL_MS: u64 = 5;
|
||||
|
||||
pub struct VoiceHeader {
|
||||
pub packet_id: u16,
|
||||
pub _flags: u16,
|
||||
pub voice_seq: u32,
|
||||
pub timestamp: u32,
|
||||
pub channel_id: u64,
|
||||
}
|
||||
|
||||
impl VoiceHeader {
|
||||
pub fn parse(buf: &[u8]) -> Option<Self> {
|
||||
if buf.len() < VOICE_HDR_SIZE {
|
||||
return None;
|
||||
}
|
||||
Some(Self {
|
||||
packet_id: u16::from_be_bytes([buf[0], buf[1]]),
|
||||
_flags: u16::from_be_bytes([buf[2], buf[3]]),
|
||||
voice_seq: u32::from_be_bytes([buf[4], buf[5], buf[6], buf[7]]),
|
||||
timestamp: u32::from_be_bytes([buf[8], buf[9], buf[10], buf[11]]),
|
||||
channel_id: u64::from_be_bytes([
|
||||
buf[12], buf[13], buf[14], buf[15], buf[16], buf[17], buf[18], buf[19],
|
||||
]),
|
||||
})
|
||||
}
|
||||
}
|
||||
12
voice-node/src/main.rs
Normal file
12
voice-node/src/main.rs
Normal file
|
|
@ -0,0 +1,12 @@
|
|||
use anyhow::Result;
|
||||
use vnox_voice_node::{load_config, runner};
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<()> {
|
||||
tracing_subscriber::fmt()
|
||||
.with_env_filter(std::env::var("VNOX_LOG").unwrap_or_else(|_| "info".into()))
|
||||
.init();
|
||||
|
||||
let config = load_config()?;
|
||||
runner::run(config).await
|
||||
}
|
||||
97
voice-node/src/relay.rs
Normal file
97
voice-node/src/relay.rs
Normal file
|
|
@ -0,0 +1,97 @@
|
|||
use crate::jitter::{BufferedPacket, JitterBuffer};
|
||||
use std::collections::HashMap;
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
use tokio::net::UdpSocket;
|
||||
use tokio::sync::RwLock;
|
||||
use tracing::debug;
|
||||
|
||||
/// Per-channel state: members, jitter buffer, and sender tracking.
|
||||
pub struct ChannelState {
|
||||
/// member address → last packet time
|
||||
pub members: HashMap<SocketAddr, Instant>,
|
||||
/// reorders and smooths voice packets
|
||||
pub jitter: JitterBuffer,
|
||||
/// voice_seq → original sender address (for relay after pop)
|
||||
pub senders: HashMap<u32, SocketAddr>,
|
||||
}
|
||||
|
||||
/// Maps channel_id → per-channel state.
|
||||
pub type ChannelMap = Arc<RwLock<HashMap<u64, ChannelState>>>;
|
||||
|
||||
pub fn new_channel_map() -> ChannelMap {
|
||||
Arc::new(RwLock::new(HashMap::new()))
|
||||
}
|
||||
|
||||
/// Drop members with no packets for longer than `max_idle`.
|
||||
pub async fn cleanup_stale(channels: &ChannelMap, max_idle: Duration) {
|
||||
let cutoff = Instant::now() - max_idle;
|
||||
let mut lock = channels.write().await;
|
||||
lock.retain(|_, state| {
|
||||
state.members.retain(|_, last_seen| *last_seen >= cutoff);
|
||||
!state.members.is_empty()
|
||||
});
|
||||
}
|
||||
|
||||
/// Push a raw voice packet into the jitter buffer for `channel_id`.
|
||||
pub async fn push_packet(
|
||||
channels: &ChannelMap,
|
||||
channel_id: u64,
|
||||
sender: SocketAddr,
|
||||
voice_seq: u32,
|
||||
timestamp: u32,
|
||||
raw_data: &[u8],
|
||||
arrived_at: u64,
|
||||
) {
|
||||
let pkt = BufferedPacket {
|
||||
voice_seq,
|
||||
timestamp,
|
||||
channel_id,
|
||||
opus_data: raw_data.to_vec(),
|
||||
arrived_at,
|
||||
};
|
||||
let mut lock = channels.write().await;
|
||||
let state = lock.entry(channel_id).or_insert_with(|| ChannelState {
|
||||
members: HashMap::new(),
|
||||
jitter: JitterBuffer::new(40, true),
|
||||
senders: HashMap::new(),
|
||||
});
|
||||
state.jitter.push(pkt);
|
||||
state.senders.insert(voice_seq, sender);
|
||||
}
|
||||
|
||||
/// Pop ready packets from every channel's jitter buffer and relay them.
|
||||
pub async fn pop_and_relay(socket: &UdpSocket, channels: &ChannelMap, now_ms: u64) {
|
||||
let mut lock = channels.write().await;
|
||||
for (_channel_id, state) in lock.iter_mut() {
|
||||
while let Some(pkt) = state.jitter.pop_ready(now_ms) {
|
||||
let sender = state.senders.remove(&pkt.voice_seq).unwrap_or_else(|| {
|
||||
tracing::warn!("sender not found for seq={}", pkt.voice_seq);
|
||||
SocketAddr::from(([0, 0, 0, 0], 0))
|
||||
});
|
||||
for &addr in state.members.keys() {
|
||||
if addr != sender
|
||||
&& let Err(e) = socket.send_to(&pkt.opus_data, addr).await
|
||||
{
|
||||
debug!("relay send error to {addr}: {e}");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Register (or refresh) a sender in the channel member set.
|
||||
pub async fn touch_member(channels: &ChannelMap, channel_id: u64, addr: SocketAddr) {
|
||||
channels
|
||||
.write()
|
||||
.await
|
||||
.entry(channel_id)
|
||||
.or_insert_with(|| ChannelState {
|
||||
members: HashMap::new(),
|
||||
jitter: JitterBuffer::new(40, true),
|
||||
senders: HashMap::new(),
|
||||
})
|
||||
.members
|
||||
.insert(addr, Instant::now());
|
||||
}
|
||||
97
voice-node/src/runner.rs
Normal file
97
voice-node/src/runner.rs
Normal file
|
|
@ -0,0 +1,97 @@
|
|||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
use tokio::net::UdpSocket;
|
||||
use tracing::{debug, error, info, warn};
|
||||
|
||||
use crate::relay;
|
||||
use crate::{
|
||||
Config, MAX_UDP_PACKET, PLAYOUT_INTERVAL_MS, VOICE_HDR_SIZE, VOICE_PACKET_ID, VoiceHeader,
|
||||
};
|
||||
|
||||
pub async fn run(config: Config) -> anyhow::Result<()> {
|
||||
run_bind(&config.node.name, &config.voice.bind).await
|
||||
}
|
||||
|
||||
pub async fn run_bind(node_name: &str, bind: &str) -> anyhow::Result<()> {
|
||||
info!("VNOX Voice Node starting — node: {node_name}");
|
||||
info!("UDP bind: {bind}");
|
||||
|
||||
let socket = UdpSocket::bind(bind).await?;
|
||||
info!("voice node listening on {bind}");
|
||||
|
||||
let channels = relay::new_channel_map();
|
||||
|
||||
let channels_cleanup = channels.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut interval = tokio::time::interval(Duration::from_secs(15));
|
||||
loop {
|
||||
interval.tick().await;
|
||||
relay::cleanup_stale(&channels_cleanup, Duration::from_secs(30)).await;
|
||||
}
|
||||
});
|
||||
|
||||
let socket = Arc::new(socket);
|
||||
let socket_relay = socket.clone();
|
||||
let channels_playout = channels.clone();
|
||||
let epoch = Instant::now();
|
||||
tokio::spawn(async move {
|
||||
let mut interval = tokio::time::interval(Duration::from_millis(PLAYOUT_INTERVAL_MS));
|
||||
loop {
|
||||
interval.tick().await;
|
||||
let now_ms = epoch.elapsed().as_millis() as u64;
|
||||
relay::pop_and_relay(socket_relay.as_ref(), &channels_playout, now_ms).await;
|
||||
}
|
||||
});
|
||||
|
||||
let mut buf = vec![0u8; MAX_UDP_PACKET];
|
||||
|
||||
loop {
|
||||
let (len, src) = match socket.recv_from(&mut buf).await {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
error!("UDP recv error: {e}");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
let data = &buf[..len];
|
||||
|
||||
let hdr = match VoiceHeader::parse(data) {
|
||||
Some(h) => h,
|
||||
None => {
|
||||
warn!("short packet from {src} ({len} bytes), dropping");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
if hdr.packet_id != VOICE_PACKET_ID {
|
||||
debug!(
|
||||
"non-voice packet 0x{:04X} from {src}, dropping",
|
||||
hdr.packet_id
|
||||
);
|
||||
continue;
|
||||
}
|
||||
|
||||
relay::touch_member(&channels, hdr.channel_id, src).await;
|
||||
|
||||
let arrived_at = epoch.elapsed().as_millis() as u64;
|
||||
|
||||
debug!(
|
||||
"voice seq={} ch={} from {src} ({} bytes opus)",
|
||||
hdr.voice_seq,
|
||||
hdr.channel_id,
|
||||
len - VOICE_HDR_SIZE,
|
||||
);
|
||||
|
||||
relay::push_packet(
|
||||
&channels,
|
||||
hdr.channel_id,
|
||||
src,
|
||||
hdr.voice_seq,
|
||||
hdr.timestamp,
|
||||
data,
|
||||
arrived_at,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue