use std::sync::{Arc, Mutex};
use std::time::Duration;
use moq_audio::playback::Input;
use tokio::time::Instant;
use super::output::{Output, Sink, Speaker};
use super::window::Event;
#[derive(Clone, Copy, Debug)]
pub(super) struct Played {
pub sink: usize,
pub sample: f32,
pub from: Instant,
pub to: Instant,
}
#[derive(Clone, Default)]
pub(super) struct Recorder {
state: Arc<Mutex<State>>,
}
#[derive(Default)]
struct State {
played: Vec<Played>,
sinks: usize,
events: Vec<Event>,
}
impl Recorder {
pub fn present(&self, media: &super::media::Media<Self>) -> Option<hang::moq_net::Timestamp> {
let mut video = media.video.lock().unwrap();
let presentation = media.presentation.lock().unwrap();
let now = Instant::now().into_std();
let mut shown = None;
while video
.front()
.is_some_and(|frame| presentation.due(frame.timestamp).is_none_or(|at| at <= now))
{
shown = video.pop_front().map(|frame| frame.timestamp);
}
if shown.is_some() {
media.drained.notify_one();
}
shown
}
pub fn played(&self) -> Vec<Played> {
self.state.lock().unwrap().played.clone()
}
pub fn events(&self) -> Vec<Event> {
std::mem::take(&mut self.state.lock().unwrap().events)
}
}
impl Output for Recorder {
type Speaker = Self;
async fn speaker(&self) -> anyhow::Result<Self> {
Ok(self.clone())
}
fn send(&self, event: Event) {
self.state.lock().unwrap().events.push(event);
}
}
impl Speaker for Recorder {
type Sink = FakeSink;
fn sink(&self, input: Input) -> anyhow::Result<FakeSink> {
anyhow::ensure!(input.format == moq_audio::Format::F32, "the fake sink only reads f32");
let id = {
let mut state = self.state.lock().unwrap();
state.sinks += 1;
state.sinks - 1
};
Ok(FakeSink {
state: self.state.clone(),
id,
stride: input.layout.channels() as usize * size_of::<f32>(),
sample_rate: input.sample_rate,
latency: input.latency,
end: Instant::now() + input.latency,
})
}
}
pub(super) struct FakeSink {
state: Arc<Mutex<State>>,
id: usize,
stride: usize,
sample_rate: u32,
latency: Duration,
end: Instant,
}
impl Sink for FakeSink {
fn write(&mut self, samples: &[u8]) -> anyhow::Result<()> {
anyhow::ensure!(samples.len().is_multiple_of(self.stride), "misaligned write");
let now = Instant::now();
if self.buffered() <= self.latency / 4 {
self.end = now + self.latency;
}
let frames = (samples.len() / self.stride) as u64;
let duration = Duration::from_nanos(frames * 1_000_000_000 / self.sample_rate as u64);
let sample = samples
.first_chunk::<4>()
.map(|bytes| f32::from_le_bytes(*bytes))
.unwrap_or_default();
self.state.lock().unwrap().played.push(Played {
sink: self.id,
sample,
from: self.end,
to: self.end + duration,
});
self.end += duration;
Ok(())
}
fn buffered(&self) -> Duration {
self.end.saturating_duration_since(Instant::now())
}
}
impl Drop for FakeSink {
fn drop(&mut self) {
let now = Instant::now();
let mut state = self.state.lock().unwrap();
state.played.retain_mut(|played| {
if played.sink == self.id {
played.to = played.to.min(now);
}
played.from < played.to
});
}
}