281 lines
8.6 KiB
Rust
281 lines
8.6 KiB
Rust
use std::collections::VecDeque;
|
|
|
|
use crate::error::{CoreError, Result};
|
|
use crate::input::decode_input;
|
|
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()
|
|
}
|
|
|
|
/// Validates and enqueues one VGI1 envelope without blocking.
|
|
///
|
|
/// # Errors
|
|
///
|
|
/// Returns a stable VGI1 parse error, `queue_full`, or `cancelled`.
|
|
pub fn enqueue_input(&mut self, bytes: Vec<u8>, features: &[&str]) -> Result<()> {
|
|
if self.cancelled {
|
|
return Err(CoreError::Cancelled);
|
|
}
|
|
decode_input(&bytes, features)?;
|
|
if self.input.len() == INPUT_QUEUE_CAPACITY {
|
|
return Err(CoreError::QueueFull);
|
|
}
|
|
self.input.push_back(bytes);
|
|
Ok(())
|
|
}
|
|
|
|
/// Removes the oldest queued VGI1 envelope.
|
|
#[must_use]
|
|
pub fn dequeue_input(&mut self) -> Option<Vec<u8>> {
|
|
self.input.pop_front()
|
|
}
|
|
|
|
/// 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(())
|
|
}
|
|
|
|
/// Removes the oldest queued reliable control body.
|
|
#[must_use]
|
|
pub fn dequeue_control(&mut self) -> Option<Vec<u8>> {
|
|
self.control.pop_front()
|
|
}
|
|
|
|
/// 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::Magic)
|
|
);
|
|
let input = b"VGI1\x01\x04\x01\x00\x00\x1e".to_vec();
|
|
for _ in 0..INPUT_QUEUE_CAPACITY {
|
|
session
|
|
.enqueue_input(input.clone(), &[])
|
|
.expect("within bound");
|
|
}
|
|
assert_eq!(session.enqueue_input(input, &[]), 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]));
|
|
}
|
|
|
|
#[test]
|
|
fn input_queue_rejects_malformed_vgi1_at_the_boundary() {
|
|
let mut session = SessionCore::new();
|
|
assert_eq!(
|
|
session.enqueue_input(vec![0], &[]),
|
|
Err(CoreError::Truncated)
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn input_queue_honors_negotiated_vgi1_features() {
|
|
let mut session = SessionCore::new();
|
|
let absolute = b"VGI1\x06\x08\x00\x01\x00\x01\x00\x02\x00\x02".to_vec();
|
|
assert_eq!(
|
|
session.enqueue_input(absolute.clone(), &[]),
|
|
Err(CoreError::UnsupportedFeature)
|
|
);
|
|
session
|
|
.enqueue_input(absolute.clone(), &["input.absolute.v1"])
|
|
.expect("negotiated absolute input");
|
|
assert_eq!(session.dequeue_input(), Some(absolute));
|
|
}
|
|
|
|
#[test]
|
|
fn input_and_control_queues_drain_in_order() {
|
|
let mut session = SessionCore::new();
|
|
let first = b"VGI1\x01\x04\x01\x00\x00\x1e".to_vec();
|
|
let second = b"VGI1\x01\x04\x00\x00\x00\x1e".to_vec();
|
|
session
|
|
.enqueue_input(first.clone(), &[])
|
|
.expect("valid input");
|
|
session
|
|
.enqueue_input(second.clone(), &[])
|
|
.expect("valid input");
|
|
session.enqueue_control(vec![1]).expect("valid control");
|
|
session.enqueue_control(vec![2]).expect("valid control");
|
|
|
|
assert_eq!(session.dequeue_input(), Some(first));
|
|
assert_eq!(session.dequeue_input(), Some(second));
|
|
assert_eq!(session.dequeue_input(), None);
|
|
assert_eq!(session.dequeue_control(), Some(vec![1]));
|
|
assert_eq!(session.dequeue_control(), Some(vec![2]));
|
|
assert_eq!(session.dequeue_control(), None);
|
|
}
|
|
}
|