feat(core): add safe Rust wire session core

This commit is contained in:
sechmachine
2026-08-12 14:03:53 +07:00
parent 12ad2a4daa
commit 48ee082c0b
15 changed files with 2115 additions and 0 deletions
+220
View File
@@ -0,0 +1,220 @@
use std::collections::VecDeque;
use crate::error::{CoreError, Result};
use crate::media::EncodedUnit;
pub const INPUT_QUEUE_CAPACITY: usize = 64;
pub const CONTROL_QUEUE_CAPACITY: usize = 64;
const MEDIA_QUEUE_CAPACITY: usize = 4;
const MAX_CONTROL_BYTES: usize = 128 * 1024;
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct SessionStats {
pub dropped_media_units: u64,
}
#[derive(Debug, Default)]
pub struct SessionCore {
input: VecDeque<Vec<u8>>,
control: VecDeque<Vec<u8>>,
media: VecDeque<EncodedUnit>,
cancelled: bool,
stats: SessionStats,
}
impl SessionCore {
#[must_use]
pub fn new() -> Self {
Self::default()
}
/// Enqueues one already-validated VGI1 envelope without blocking.
///
/// # Errors
///
/// Returns `invalid_argument`, `queue_full`, or `cancelled`.
pub fn enqueue_input(&mut self, bytes: Vec<u8>) -> Result<()> {
if self.cancelled {
return Err(CoreError::Cancelled);
}
if bytes.is_empty() || bytes.len() > 23 {
return Err(CoreError::InvalidArgument);
}
if self.input.len() == INPUT_QUEUE_CAPACITY {
return Err(CoreError::QueueFull);
}
self.input.push_back(bytes);
Ok(())
}
/// Enqueues one bounded reliable control body without blocking.
///
/// # Errors
///
/// Returns `invalid_argument`, `queue_full`, or `cancelled`.
pub fn enqueue_control(&mut self, bytes: Vec<u8>) -> Result<()> {
if self.cancelled {
return Err(CoreError::Cancelled);
}
if bytes.len() > MAX_CONTROL_BYTES {
return Err(CoreError::InvalidArgument);
}
if self.control.len() == CONTROL_QUEUE_CAPACITY {
return Err(CoreError::QueueFull);
}
self.control.push_back(bytes);
Ok(())
}
/// Enqueues one bounded complete encoded media unit.
///
/// # Errors
///
/// Returns `invalid_argument` for an oversized unit or `cancelled` after cancellation.
pub fn enqueue_media(&mut self, unit: EncodedUnit) -> Result<()> {
if unit.payload.len() > crate::media::MAX_COMPLETE_UNIT_BYTES {
return Err(CoreError::InvalidArgument);
}
if self.cancelled {
self.stats.dropped_media_units += 1;
return Err(CoreError::Cancelled);
}
if self.media.len() == MEDIA_QUEUE_CAPACITY {
let position = self
.media
.iter()
.position(|queued| queued.channel == unit.channel);
if let Some(position) = position {
self.media.remove(position);
} else {
self.stats.dropped_media_units += 1;
return Ok(());
}
self.stats.dropped_media_units += 1;
}
self.media.push_back(unit);
Ok(())
}
pub fn pop_media(&mut self) -> Option<EncodedUnit> {
self.media.pop_front()
}
pub fn cancel(&mut self) {
self.cancelled = true;
self.input.clear();
self.control.clear();
self.media.clear();
}
#[must_use]
pub const fn is_cancelled(&self) -> bool {
self.cancelled
}
#[must_use]
pub const fn stats(&self) -> SessionStats {
self.stats
}
}
#[cfg(test)]
#[derive(Debug, Default)]
struct InProcessSession {
incoming: VecDeque<Vec<u8>>,
}
#[cfg(test)]
impl InProcessSession {
fn push_incoming(&mut self, bytes: Vec<u8>) {
self.incoming.push_back(bytes);
}
fn pop_incoming(&mut self) -> Option<Vec<u8>> {
self.incoming.pop_front()
}
}
#[cfg(test)]
mod tests {
use super::{InProcessSession, SessionCore, CONTROL_QUEUE_CAPACITY, INPUT_QUEUE_CAPACITY};
use crate::error::CoreError;
use crate::media::{EncodedUnit, MediaChannel};
#[test]
fn queues_are_bounded_and_cancellation_is_idempotent() {
let mut session = SessionCore::new();
assert_eq!(
session.enqueue_input(vec![0; 24]),
Err(CoreError::InvalidArgument)
);
for value in 0..INPUT_QUEUE_CAPACITY {
session
.enqueue_input(vec![u8::try_from(value).expect("capacity fits u8")])
.expect("within bound");
}
assert_eq!(session.enqueue_input(vec![0]), Err(CoreError::QueueFull));
for value in 0..CONTROL_QUEUE_CAPACITY {
session
.enqueue_control(vec![u8::try_from(value).expect("capacity fits u8")])
.expect("within bound");
}
assert_eq!(session.enqueue_control(vec![0]), Err(CoreError::QueueFull));
let mut oversized_control = SessionCore::new();
assert_eq!(
oversized_control.enqueue_control(vec![0; 128 * 1024 + 1]),
Err(CoreError::InvalidArgument)
);
session
.enqueue_media(EncodedUnit {
channel: MediaChannel::Audio,
sequence: 1,
timestamp_ms: 1,
payload: vec![1],
})
.expect("bounded media");
session.cancel();
session.cancel();
assert!(session.is_cancelled());
assert_eq!(session.enqueue_input(vec![0]), Err(CoreError::Cancelled));
assert!(session.pop_media().is_none());
}
#[test]
fn media_queue_evicts_oldest_same_channel_at_four_units() {
let mut session = SessionCore::new();
assert_eq!(
session.enqueue_media(EncodedUnit {
channel: MediaChannel::Video,
sequence: 99,
timestamp_ms: 0,
payload: vec![0; 1_048_577],
}),
Err(CoreError::InvalidArgument)
);
for sequence in 0..5 {
session
.enqueue_media(EncodedUnit {
channel: MediaChannel::Video,
sequence,
timestamp_ms: u64::from(sequence),
payload: vec![u8::try_from(sequence).expect("test sequence fits u8")],
})
.expect("bounded media");
}
assert_eq!(session.stats().dropped_media_units, 1);
assert_eq!(session.pop_media().expect("media").sequence, 1);
}
#[test]
fn private_in_process_session_preserves_byte_order() {
let mut transport = InProcessSession::default();
transport.push_incoming(vec![1, 2]);
transport.push_incoming(vec![3]);
assert_eq!(transport.pop_incoming(), Some(vec![1, 2]));
assert_eq!(transport.pop_incoming(), Some(vec![3]));
}
}