use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use parking_lot::{Condvar, Mutex, RwLock};
use tokio::sync::watch;
pub const NUM_BARS: usize = 48;
const DEFAULT_FPS: u8 = 60;
pub const WAVEFORM_SAMPLES: usize = 2048;
const DELAY_LINE_SIZE: usize = crate::player::RING_BUFFER_SIZE + WAVEFORM_SAMPLES * 2;
#[derive(Clone)]
pub struct VizFrame {
pub spectrum: [f32; NUM_BARS],
pub peaks: [f32; NUM_BARS],
pub vu_levels: [f32; 2],
pub beat_energy: f32,
pub timestamp: std::time::Instant,
pub waveform: Vec<f32>,
}
impl Default for VizFrame {
fn default() -> Self {
Self {
spectrum: [0.0; NUM_BARS],
peaks: [0.0; NUM_BARS],
vu_levels: [0.0; 2],
beat_energy: 0.0,
timestamp: std::time::Instant::now(),
waveform: Vec::new(),
}
}
}
pub struct VizSnapshot {
inner: RwLock<VizFrame>,
reads: AtomicU64,
published: watch::Sender<u64>,
park: Mutex<bool>,
unpark: Condvar,
parked: AtomicBool,
interval_us: AtomicU64,
}
impl VizSnapshot {
pub fn new() -> Arc<Self> {
Arc::new(Self {
inner: RwLock::new(VizFrame::default()),
reads: AtomicU64::new(0),
published: watch::Sender::new(0),
park: Mutex::new(false),
unpark: Condvar::new(),
parked: AtomicBool::new(false),
interval_us: AtomicU64::new(Self::interval_us(DEFAULT_FPS)),
})
}
pub fn read(&self) -> VizFrame {
self.touch();
self.inner.read().clone()
}
pub fn touch(&self) {
self.reads.fetch_add(1, Ordering::Relaxed);
if self.parked.load(Ordering::Relaxed) {
self.wake();
}
}
pub fn wake(&self) {
let mut pending = self.park.lock();
*pending = true;
self.unpark.notify_all();
}
pub fn park_while_idle(&self, still_idle: impl Fn() -> bool) {
let mut pending = self.park.lock();
if std::mem::take(&mut pending) || !still_idle() {
return;
}
self.parked.store(true, Ordering::Relaxed);
while !*pending {
self.unpark.wait(&mut pending);
}
*pending = false;
self.parked.store(false, Ordering::Relaxed);
}
pub fn subscribe(&self) -> watch::Receiver<u64> {
self.published.subscribe()
}
pub fn reads(&self) -> u64 {
self.reads.load(Ordering::Relaxed)
}
pub fn set_fps(&self, fps: u8) {
self.interval_us
.store(Self::interval_us(fps), Ordering::Relaxed);
}
pub fn interval(&self) -> std::time::Duration {
std::time::Duration::from_micros(self.interval_us.load(Ordering::Relaxed))
}
pub fn fps(&self) -> u8 {
(1_000_000 / self.interval_us.load(Ordering::Relaxed).max(1)) as u8
}
fn interval_us(fps: u8) -> u64 {
1_000_000 / fps.clamp(1, 240) as u64
}
pub fn write(&self, frame: VizFrame) {
*self.inner.write() = frame;
self.published.send_modify(|version| *version += 1);
}
pub fn levels(&self) -> VizLevels {
self.touch();
let frame = self.inner.read();
VizLevels::of(&frame.spectrum)
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq)]
pub struct VizLevels {
pub low: f32,
pub mid: f32,
pub high: f32,
}
impl VizLevels {
fn of(spectrum: &[f32; NUM_BARS]) -> Self {
let band = NUM_BARS / 3;
let mean = |bars: &[f32]| bars.iter().sum::<f32>() / bars.len() as f32;
Self {
low: mean(&spectrum[..band]),
mid: mean(&spectrum[band..band * 2]),
high: mean(&spectrum[band * 2..]),
}
}
}
#[derive(Default)]
pub struct RawVizSnapshot {
pub samples: Vec<f32>,
pub channels: u16,
pub sample_rate: u32,
}
struct VizSamples {
buf: Vec<f32>,
write_pos: usize,
head_offset: u64,
channels: u16,
sample_rate: u32,
}
pub struct VizBuffer {
samples: Mutex<VizSamples>,
}
impl VizBuffer {
pub fn new() -> Arc<Self> {
Arc::new(Self {
samples: Mutex::new(VizSamples {
buf: vec![0.0; DELAY_LINE_SIZE],
write_pos: 0,
head_offset: 0,
channels: 2,
sample_rate: 44100,
}),
})
}
pub fn push_samples(&self, samples: &[f32], channels: u16, sample_rate: u32) {
let mut inner = self.samples.lock();
inner.channels = channels;
inner.sample_rate = sample_rate;
inner.head_offset += samples.len() as u64;
let buf_len = inner.buf.len();
if samples.len() >= buf_len {
let start = samples.len() - buf_len;
inner.buf.copy_from_slice(&samples[start..]);
inner.write_pos = 0;
} else {
let pos = inner.write_pos;
let first = buf_len - pos;
if samples.len() <= first {
inner.buf[pos..pos + samples.len()].copy_from_slice(samples);
inner.write_pos = (pos + samples.len()) % buf_len;
} else {
inner.buf[pos..].copy_from_slice(&samples[..first]);
let remaining = samples.len() - first;
inner.buf[..remaining].copy_from_slice(&samples[first..]);
inner.write_pos = remaining;
}
}
}
pub fn reset(&self) {
let mut inner = self.samples.lock();
inner.buf.fill(0.0);
inner.write_pos = 0;
inner.head_offset = 0;
}
pub fn snapshot_at(&self, played: u64, frames: usize, out: &mut RawVizSnapshot) {
let inner = self.samples.lock();
let buf_len = inner.buf.len();
out.channels = inner.channels;
out.sample_rate = inner.sample_rate;
out.samples.clear();
let wanted = frames
.saturating_mul(inner.channels.max(1) as usize)
.min(buf_len);
if wanted == 0 {
return;
}
let delay = inner
.head_offset
.saturating_sub(played)
.min((buf_len - wanted) as u64) as usize;
let end = (inner.write_pos + buf_len - delay) % buf_len;
let start = (end + buf_len - wanted) % buf_len;
out.samples.reserve(wanted);
if start < end {
out.samples.extend_from_slice(&inner.buf[start..end]);
} else {
out.samples.extend_from_slice(&inner.buf[start..]);
out.samples.extend_from_slice(&inner.buf[..end]);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn push_ramp(buf: &VizBuffer, count: usize) {
let samples: Vec<f32> = (0..count).map(|i| i as f32).collect();
buf.push_samples(&samples, 2, 44100);
}
#[test]
fn snapshot_at_head_returns_newest_samples() {
let buf = VizBuffer::new();
push_ramp(&buf, 1000);
let mut snap = RawVizSnapshot::default();
buf.snapshot_at(1000, 100, &mut snap);
assert_eq!(snap.samples.len(), 200);
for (i, &val) in snap.samples.iter().enumerate() {
assert_eq!(val, (800 + i) as f32);
}
}
#[test]
fn snapshot_at_walks_back_to_the_play_head() {
let buf = VizBuffer::new();
push_ramp(&buf, 100_000);
let mut snap = RawVizSnapshot::default();
buf.snapshot_at(60_000, 512, &mut snap);
assert_eq!(snap.samples.len(), 1024);
assert_eq!(*snap.samples.last().unwrap(), 59_999.0);
assert_eq!(snap.samples[0], (60_000 - 1024) as f32);
}
#[test]
fn snapshot_at_tracks_the_play_head_across_wraps() {
let buf = VizBuffer::new();
let total = DELAY_LINE_SIZE + 5_000;
push_ramp(&buf, total);
let played = (total - 2_000) as u64;
let mut snap = RawVizSnapshot::default();
buf.snapshot_at(played, 256, &mut snap);
assert_eq!(snap.samples.len(), 512);
assert_eq!(*snap.samples.last().unwrap(), (played - 1) as f32);
assert_eq!(snap.samples[0], (played - 512) as f32);
}
#[test]
fn snapshot_at_clamps_lookback_to_buffer_length() {
let buf = VizBuffer::new();
push_ramp(&buf, DELAY_LINE_SIZE * 2);
let mut snap = RawVizSnapshot::default();
buf.snapshot_at(0, 64, &mut snap);
assert_eq!(snap.samples.len(), 128);
assert_eq!(snap.samples[0], DELAY_LINE_SIZE as f32);
}
#[test]
fn reset_clears_samples_and_offset() {
let buf = VizBuffer::new();
push_ramp(&buf, 10_000);
buf.reset();
let mut snap = RawVizSnapshot::default();
buf.snapshot_at(0, 32, &mut snap);
assert!(snap.samples.iter().all(|&s| s == 0.0));
push_ramp(&buf, 500);
buf.snapshot_at(500, 10, &mut snap);
assert_eq!(*snap.samples.last().unwrap(), 499.0);
}
#[test]
fn snapshot_at_reports_metadata() {
let buf = VizBuffer::new();
buf.push_samples(&[1.0, 2.0], 1, 96000);
let mut snap = RawVizSnapshot::default();
buf.snapshot_at(2, 2, &mut snap);
assert_eq!(snap.channels, 1);
assert_eq!(snap.sample_rate, 96000);
assert_eq!(snap.samples, vec![1.0, 2.0]);
}
#[test]
fn fps_sets_the_analysis_interval_and_clamps() {
let snap = VizSnapshot::new();
assert_eq!(snap.fps(), DEFAULT_FPS);
snap.set_fps(120);
assert_eq!(snap.fps(), 120);
assert_eq!(snap.interval(), std::time::Duration::from_micros(8_333));
snap.set_fps(0);
assert_eq!(snap.fps(), 1);
}
#[test]
fn a_wake_landing_before_the_park_is_not_slept_through() {
let snap = VizSnapshot::new();
snap.wake();
snap.park_while_idle(|| true);
let woken = Arc::new(AtomicBool::new(false));
let flag = Arc::clone(&woken);
let snapshot = Arc::clone(&snap);
let waiter = std::thread::spawn(move || {
snapshot.park_while_idle(|| true);
flag.store(true, Ordering::Relaxed);
});
std::thread::sleep(std::time::Duration::from_millis(50));
assert!(
!woken.load(Ordering::Relaxed),
"parked thread woke on its own"
);
snap.wake();
waiter.join().unwrap();
}
#[test]
fn a_read_wakes_a_parked_analyser() {
let snap = VizSnapshot::new();
let snapshot = Arc::clone(&snap);
let waiter = std::thread::spawn(move || snapshot.park_while_idle(|| true));
std::thread::sleep(std::time::Duration::from_millis(50));
snap.touch();
waiter.join().unwrap();
}
#[test]
fn viz_snapshot_read_write() {
let snap = VizSnapshot::new();
let frame = snap.read();
assert_eq!(frame.spectrum.len(), NUM_BARS);
assert_eq!(frame.vu_levels, [0.0, 0.0]);
let mut new_spectrum = [0.0f32; NUM_BARS];
new_spectrum[5] = 0.9;
snap.write(VizFrame {
spectrum: new_spectrum,
peaks: [0.0; NUM_BARS],
vu_levels: [0.5, 0.5],
beat_energy: 0.0,
timestamp: std::time::Instant::now(),
waveform: Vec::new(),
});
let frame2 = snap.read();
assert!((frame2.spectrum[5] - 0.9).abs() < 0.001);
assert!((frame2.vu_levels[0] - 0.5).abs() < 0.001);
}
#[test]
fn levels_average_each_third_of_the_bars() {
let mut spectrum = [0.0f32; NUM_BARS];
let band = NUM_BARS / 3;
spectrum[..band].fill(0.6);
spectrum[band..band * 2].fill(0.3);
spectrum[band * 2..].fill(0.0);
spectrum[NUM_BARS - 1] = 1.0;
let levels = VizLevels::of(&spectrum);
assert!((levels.low - 0.6).abs() < 0.001);
assert!((levels.mid - 0.3).abs() < 0.001);
assert!((levels.high - 1.0 / band as f32).abs() < 0.001);
}
}