use std::collections::BTreeMap;
use sim_kernel::{Result, Symbol};
use sim_lib_audio_graph_core::Transport;
use sim_lib_stream_clock::Clock;
use sim_lib_stream_core::{StreamDiagnostic, StreamEnvelope, StreamMedia, StreamPacket};
pub fn live_clock_symbol() -> Symbol {
Symbol::qualified("clock", "audio-graph-live")
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct LiveTransportClock {
clock: Clock,
sample_rate_hz: u32,
}
#[derive(Clone, Debug)]
pub struct LanBufferedPreviewWindow {
next_sequence: u64,
reorder_depth: u64,
pending: BTreeMap<u64, StreamEnvelope>,
diagnostics: Vec<StreamPacket>,
}
impl LiveTransportClock {
pub fn sample_frame(sample_rate_hz: u32) -> Result<Self> {
Ok(Self {
clock: Clock::frame(live_clock_symbol(), u64::from(sample_rate_hz))?,
sample_rate_hz,
})
}
pub fn clock(&self) -> &Clock {
&self.clock
}
pub fn sample_rate_hz(&self) -> u32 {
self.sample_rate_hz
}
pub fn transport_at(&self, sample_pos: u64, playing: bool) -> Transport {
Transport {
playing,
sample_pos,
tempo_bpm: 120.0,
ppq_pos: 0.0,
}
}
}
impl LanBufferedPreviewWindow {
pub fn new(next_sequence: u64, reorder_depth: u64) -> Self {
Self {
next_sequence,
reorder_depth,
pending: BTreeMap::new(),
diagnostics: Vec::new(),
}
}
pub fn push(&mut self, envelope: StreamEnvelope) -> Result<Vec<StreamEnvelope>> {
validate_lan_buffered_preview_envelope(&envelope)?;
let sequence = envelope.sequence();
if sequence < self.next_sequence || self.pending.contains_key(&sequence) {
self.record(
lan_buffered_preview_late_packet_diagnostic_kind(),
format!(
"LAN preview packet {sequence} arrived after sequence {} was expected",
self.next_sequence
),
);
return Ok(Vec::new());
}
if sequence > self.next_sequence {
let gap = sequence - self.next_sequence;
self.record(
lan_buffered_preview_jitter_diagnostic_kind(),
format!("LAN preview packet {sequence} arrived with a sequence gap of {gap}"),
);
self.record(
lan_buffered_preview_reorder_diagnostic_kind(),
format!(
"LAN preview packet {sequence} arrived before expected sequence {}",
self.next_sequence
),
);
if gap > self.reorder_depth {
self.record(
lan_buffered_preview_drop_diagnostic_kind(),
format!(
"LAN preview dropped missing sequence range {}..{sequence}",
self.next_sequence
),
);
self.next_sequence = sequence;
return Ok(self.accept_ready(envelope));
}
self.pending.insert(sequence, envelope);
return Ok(Vec::new());
}
Ok(self.accept_ready(envelope))
}
pub fn drain_diagnostics(&mut self) -> Vec<StreamPacket> {
std::mem::take(&mut self.diagnostics)
}
fn accept_ready(&mut self, envelope: StreamEnvelope) -> Vec<StreamEnvelope> {
let mut ready = vec![envelope];
self.next_sequence = self.next_sequence.saturating_add(1);
while let Some(envelope) = self.pending.remove(&self.next_sequence) {
ready.push(envelope);
self.next_sequence = self.next_sequence.saturating_add(1);
}
ready
}
fn record(&mut self, kind: Symbol, message: String) {
self.diagnostics
.push(StreamPacket::Diagnostic(StreamDiagnostic::new(
kind, message,
)));
}
}
pub fn validate_lan_buffered_preview_envelope(envelope: &StreamEnvelope) -> Result<()> {
if envelope.media() != StreamMedia::Pcm {
return Err(sim_kernel::Error::Eval(
"LAN buffered audio preview requires PCM media".to_owned(),
));
}
if envelope.profile().name()
!= sim_lib_stream_core::TransportProfile::lan_buffered_audio_preview().name()
{
return Err(sim_kernel::Error::Eval(
"LAN buffered audio preview requires stream/profile/lan-buffered-audio-preview"
.to_owned(),
));
}
Ok(())
}
pub fn lan_buffered_preview_jitter_diagnostic_kind() -> Symbol {
Symbol::qualified("stream/preview", "Jitter")
}
pub fn lan_buffered_preview_drop_diagnostic_kind() -> Symbol {
Symbol::qualified("stream/preview", "Drop")
}
pub fn lan_buffered_preview_reorder_diagnostic_kind() -> Symbol {
Symbol::qualified("stream/preview", "Reorder")
}
pub fn lan_buffered_preview_late_packet_diagnostic_kind() -> Symbol {
Symbol::qualified("stream/preview", "LatePacket")
}