use std::collections::VecDeque;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, LazyLock, Weak};
use std::time::{Duration, Instant};
use parking_lot::{Condvar, Mutex};
use serde::{Deserialize, Serialize};
use crate::audio::viz::{VizLevels, VizSnapshot};
use crate::remote::link::LinkCommand;
use crate::remote::wire::Waker;
use crate::signal::Wake;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct Frame(pub u64, pub u16, pub u16, pub u16);
impl Frame {
pub fn new(at_ms: u64, levels: VizLevels) -> Self {
let q = |v: f32| (v.clamp(0.0, 1.0) * 1000.0).round() as u16;
Self(at_ms, q(levels.low), q(levels.mid), q(levels.high))
}
pub fn at_ms(&self) -> u64 {
self.0
}
pub fn levels(&self) -> VizLevels {
VizLevels {
low: f32::from(self.1) / 1000.0,
mid: f32::from(self.2) / 1000.0,
high: f32::from(self.3) / 1000.0,
}
}
}
pub struct Feed {
inner: Mutex<FeedState>,
changed: Condvar,
next_id: AtomicU64,
}
struct FeedState {
source: Option<Arc<Source>>,
watchers: Vec<(u64, Weak<Waker>)>,
latest: Option<(u64, Frame)>,
pumping: bool,
}
struct Source {
viz: Arc<VizSnapshot>,
position_ms: Box<dyn Fn() -> u64 + Send + Sync>,
}
static FEED: LazyLock<Arc<Feed>> = LazyLock::new(Feed::new);
pub fn feed() -> &'static Arc<Feed> {
&FEED
}
impl Feed {
pub fn new() -> Arc<Self> {
Arc::new(Self {
inner: Mutex::new(FeedState {
source: None,
watchers: Vec::new(),
latest: None,
pumping: false,
}),
changed: Condvar::new(),
next_id: AtomicU64::new(1),
})
}
pub fn provide(
self: &Arc<Self>,
viz: Arc<VizSnapshot>,
position_ms: impl Fn() -> u64 + Send + Sync + 'static,
) {
let mut inner = self.inner.lock();
inner.source = Some(Arc::new(Source {
viz,
position_ms: Box::new(position_ms),
}));
self.changed.notify_all();
}
pub fn watch(self: &Arc<Self>, waker: &Arc<Waker>) -> Watch {
let id = self.next_id.fetch_add(1, Ordering::Relaxed);
let mut inner = self.inner.lock();
inner.watchers.push((id, Arc::downgrade(waker)));
let sent = inner.latest.map_or(0, |(seq, _)| seq);
if !inner.pumping {
inner.pumping = true;
let feed = Arc::clone(self);
std::thread::Builder::new()
.name("koan-levels".into())
.spawn(move || feed.pump())
.expect("failed to spawn the levels thread");
}
self.changed.notify_all();
Watch {
feed: Arc::downgrade(self),
id,
sent,
}
}
pub fn watchers(&self) -> usize {
self.inner.lock().watchers.len()
}
fn pump(&self) {
loop {
let source = {
let mut inner = self.inner.lock();
loop {
inner.watchers.retain(|(_, w)| w.strong_count() > 0);
match &inner.source {
Some(source) if !inner.watchers.is_empty() => break Arc::clone(source),
_ => self.changed.wait(&mut inner),
}
}
};
let seen = source.viz.frames().generation();
source.viz.touch();
source.viz.frames().wait(seen);
let frame = Frame::new((source.position_ms)(), source.viz.levels());
let mut inner = self.inner.lock();
let seq = inner.latest.map_or(0, |(seq, _)| seq) + 1;
inner.latest = Some((seq, frame));
inner.watchers.retain(|(_, w)| match w.upgrade() {
Some(w) => {
w.wake();
true
}
None => false,
});
}
}
#[cfg(test)]
fn publish_for_test(&self, frame: Frame) {
let mut inner = self.inner.lock();
let seq = inner.latest.map_or(0, |(seq, _)| seq) + 1;
inner.latest = Some((seq, frame));
}
}
pub struct Watch {
feed: Weak<Feed>,
id: u64,
sent: u64,
}
impl Watch {
pub fn take(&mut self) -> Option<Frame> {
let feed = self.feed.upgrade()?;
let (seq, frame) = feed.inner.lock().latest?;
(seq > self.sent).then(|| {
self.sent = seq;
frame
})
}
}
impl Drop for Watch {
fn drop(&mut self) {
if let Some(feed) = self.feed.upgrade() {
feed.inner.lock().watchers.retain(|(id, _)| *id != self.id);
}
}
}
const DELAY_FRAMES: f32 = 2.0;
const JUMP_MS: u64 = 1_000;
const BACK_MS: u64 = 50;
const EASE_HALF_LIFE: Duration = Duration::from_millis(80);
const REST: f32 = 0.001;
const MAX_FRAMES: usize = 240;
#[derive(Debug)]
pub struct Interp {
frames: VecDeque<(u64, VizLevels)>,
spacing_ms: f32,
shown: VizLevels,
sampled: Option<Instant>,
}
impl Default for Interp {
fn default() -> Self {
Self {
frames: VecDeque::new(),
spacing_ms: 1000.0 / 60.0,
shown: VizLevels::default(),
sampled: None,
}
}
}
impl Interp {
pub fn push(&mut self, at_ms: u64, levels: VizLevels) {
if let Some(&(newest, _)) = self.frames.back() {
if at_ms + BACK_MS < newest || at_ms > newest + JUMP_MS {
self.frames.clear();
} else if at_ms <= newest {
if let Some(back) = self.frames.back_mut() {
*back = (newest, levels);
}
return;
} else {
let gap = (at_ms - newest) as f32;
if gap < 200.0 {
self.spacing_ms = self.spacing_ms * 0.8 + gap * 0.2;
}
}
}
self.frames.push_back((at_ms, levels));
while self.frames.len() > MAX_FRAMES {
self.frames.pop_front();
}
}
pub fn clear(&mut self) {
self.frames.clear();
}
pub fn delay_ms(&self) -> u64 {
(self.spacing_ms * DELAY_FRAMES).round() as u64
}
pub fn sample(&mut self, playhead_ms: u64, playing: bool, now: Instant) -> VizLevels {
let elapsed = self
.sampled
.map_or(Duration::ZERO, |at| now.saturating_duration_since(at));
self.sampled = Some(now);
let at = playhead_ms.saturating_sub(self.delay_ms());
match self.between(at).filter(|_| playing) {
Some(levels) => self.shown = levels,
None => {
let keep = 0.5f32.powf(elapsed.as_secs_f32() / EASE_HALF_LIFE.as_secs_f32());
self.shown = VizLevels {
low: self.shown.low * keep,
mid: self.shown.mid * keep,
high: self.shown.high * keep,
};
let stale = self
.frames
.back()
.is_some_and(|&(newest, _)| at > newest + JUMP_MS || at + JUMP_MS < newest);
if stale || !playing {
self.frames.clear();
}
}
}
self.shown
}
fn between(&mut self, at: u64) -> Option<VizLevels> {
let after = self.frames.iter().position(|&(t, _)| t >= at)?;
if after == 0 {
return (self.frames[0].0 == at).then_some(self.frames[0].1);
}
self.frames.drain(..after - 1);
let (t0, a) = self.frames[0];
let (t1, b) = self.frames[1];
let f = (at - t0) as f32 / (t1 - t0) as f32;
let lerp = |x: f32, y: f32| x + (y - x) * f;
Some(VizLevels {
low: lerp(a.low, b.low),
mid: lerp(a.mid, b.mid),
high: lerp(a.high, b.high),
})
}
pub fn moving(&self) -> bool {
!self.frames.is_empty()
|| self.shown.low > REST
|| self.shown.mid > REST
|| self.shown.high > REST
}
}
const STALL: Duration = Duration::from_secs(3);
pub struct Remote {
inner: Mutex<RemoteState>,
ticks: Wake,
tick: Condvar,
send: Box<SendWatch>,
}
type SendWatch = dyn Fn(&str, LinkCommand) -> bool + Send + Sync;
struct RemoteState {
views: Vec<Viewed>,
drawing: Option<String>,
interp: Interp,
heard: Option<Instant>,
interval: Duration,
ticking: bool,
}
struct Viewed {
target: String,
count: usize,
asked: Option<u64>,
}
static REMOTE: LazyLock<Arc<Remote>> =
LazyLock::new(|| Remote::new(crate::remote::devices::send_live));
pub fn remote() -> &'static Arc<Remote> {
&REMOTE
}
impl Remote {
pub fn new(send: impl Fn(&str, LinkCommand) -> bool + Send + Sync + 'static) -> Arc<Self> {
Arc::new(Self {
inner: Mutex::new(RemoteState {
views: Vec::new(),
drawing: None,
interp: Interp::default(),
heard: None,
interval: Duration::from_micros(1_000_000 / 60),
ticking: false,
}),
ticks: Wake::new(),
tick: Condvar::new(),
send: Box::new(send),
})
}
pub fn view(self: &Arc<Self>, target: String) -> View {
let mut inner = self.inner.lock();
match inner.views.iter_mut().find(|v| v.target == target) {
Some(v) => v.count += 1,
None => inner.views.push(Viewed {
target: target.clone(),
count: 1,
asked: None,
}),
}
if inner.drawing.as_deref() != Some(&target) {
inner.drawing = Some(target.clone());
inner.interp = Interp::default();
inner.heard = None;
}
if !inner.ticking {
inner.ticking = true;
let remote = Arc::clone(self);
std::thread::Builder::new()
.name("koan-levels-tick".into())
.spawn(move || remote.run_ticks())
.expect("failed to spawn the levels ticker");
}
self.tick.notify_all();
drop(inner);
self.ask(false);
View {
remote: Arc::downgrade(self),
target,
}
}
pub fn ask(&self, again: bool) {
let version = crate::remote::devices::version();
let due: Vec<String> = {
let inner = self.inner.lock();
inner
.views
.iter()
.filter(|v| again || v.asked != Some(version))
.map(|v| v.target.clone())
.collect()
};
for target in due {
let sent = (self.send)(&target, LinkCommand::WatchLevels { on: true });
if !sent {
log::debug!("levels: no way to {target} now; asking when one opens");
}
if let Some(v) = self
.inner
.lock()
.views
.iter_mut()
.find(|v| v.target == target)
{
v.asked = sent.then_some(version);
}
}
}
pub fn received(&self, from: &str, frame: Frame) {
let mut inner = self.inner.lock();
if inner.drawing.as_deref() != Some(from) || !inner.views.iter().any(|v| v.target == from) {
return;
}
inner.interp.push(frame.at_ms(), frame.levels());
inner.heard = Some(Instant::now());
self.tick.notify_all();
}
pub fn sample(&self, playhead_ms: u64, playing: bool) -> VizLevels {
self.inner
.lock()
.interp
.sample(playhead_ms, playing, Instant::now())
}
pub fn set_fps(&self, fps: u8) {
self.inner.lock().interval = Duration::from_micros(1_000_000 / fps.clamp(1, 240) as u64);
}
pub fn ticks(&self) -> &Wake {
&self.ticks
}
fn release(&self, target: &str) {
let mut inner = self.inner.lock();
let Some(at) = inner.views.iter().position(|v| v.target == target) else {
return;
};
inner.views[at].count -= 1;
if inner.views[at].count > 0 {
return;
}
inner.views.remove(at);
if inner.drawing.as_deref() == Some(target) {
inner.drawing = None;
inner.interp = Interp::default();
inner.heard = None;
}
drop(inner);
(self.send)(target, LinkCommand::WatchLevels { on: false });
}
fn run_ticks(&self) {
loop {
let interval = {
let mut inner = self.inner.lock();
loop {
if inner.views.is_empty() {
self.tick.wait(&mut inner);
} else if inner.interp.moving() {
break inner.interval;
} else if self.tick.wait_for(&mut inner, STALL).timed_out()
&& inner.heard.is_none_or(|at| at.elapsed() >= STALL)
&& !inner.views.is_empty()
{
drop(inner);
let playing = crate::remote::devices::target_playhead()
.is_some_and(|(_, playing)| playing);
if playing {
self.ask(true);
}
inner = self.inner.lock();
}
}
};
std::thread::sleep(interval);
self.ticks.bump();
}
}
#[cfg(test)]
fn viewers(&self, target: &str) -> usize {
self.inner
.lock()
.views
.iter()
.find(|v| v.target == target)
.map_or(0, |v| v.count)
}
}
pub struct View {
remote: Weak<Remote>,
target: String,
}
impl View {
pub fn target(&self) -> &str {
&self.target
}
}
impl Drop for View {
fn drop(&mut self) {
if let Some(remote) = self.remote.upgrade() {
remote.release(&self.target);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
type Sent = Arc<Mutex<Vec<(String, bool)>>>;
fn recording(through: bool) -> (Arc<Remote>, Sent, Arc<std::sync::atomic::AtomicBool>) {
let sent = Arc::new(Mutex::new(Vec::new()));
let open = Arc::new(std::sync::atomic::AtomicBool::new(through));
let (log, gate) = (Arc::clone(&sent), Arc::clone(&open));
let remote = Remote::new(move |to, cmd| {
if let LinkCommand::WatchLevels { on } = cmd {
log.lock().push((to.to_string(), on));
}
gate.load(Ordering::Relaxed)
});
(remote, sent, open)
}
#[test]
fn switching_devices_tells_the_old_one_to_stop() {
let (remote, sent, _) = recording(true);
let x = remote.view("x".into());
let y = remote.view("y".into());
drop(x);
assert_eq!(remote.viewers("x"), 0);
assert!(
sent.lock().contains(&("x".into(), false)),
"{:?}",
sent.lock()
);
assert!(!sent.lock().contains(&("y".into(), false)));
drop(y);
assert!(sent.lock().contains(&("y".into(), false)));
}
#[test]
fn a_device_stays_watched_while_any_view_of_it_is_held() {
let (remote, sent, _) = recording(true);
let a = remote.view("x".into());
let b = remote.view("x".into());
drop(a);
assert!(!sent.lock().contains(&("x".into(), false)));
drop(b);
assert!(sent.lock().contains(&("x".into(), false)));
}
#[test]
fn an_ask_that_found_no_way_through_is_made_again() {
let (remote, sent, open) = recording(false);
let _view = remote.view("x".into());
let asks = |sent: &Mutex<Vec<(String, bool)>>| {
sent.lock().iter().filter(|(t, on)| t == "x" && *on).count()
};
assert_eq!(asks(&sent), 1);
open.store(true, Ordering::Relaxed);
remote.ask(false);
assert_eq!(asks(&sent), 2, "asked again, having not got through");
remote.ask(true);
assert_eq!(asks(&sent), 3, "and on a renewal");
}
#[test]
fn frames_from_a_device_not_viewed_are_ignored() {
let (remote, _, _) = recording(true);
let _view = remote.view("x".into());
remote.received("y", Frame(100, 500, 500, 500));
assert!(!remote.inner.lock().interp.moving());
remote.received("x", Frame(100, 500, 500, 500));
assert!(remote.inner.lock().interp.moving());
}
#[test]
fn a_short_skip_back_clears_the_buffer() {
let mut i = interp(400, 10);
i.push(150, lv(0.9));
assert_eq!(i.frames.len(), 1);
}
fn lv(v: f32) -> VizLevels {
VizLevels {
low: v,
mid: v / 2.0,
high: v / 4.0,
}
}
fn interp(from: u64, frames: u64) -> Interp {
let mut i = Interp::default();
for n in 0..frames {
i.push(from + n * 20, lv(n as f32 / 100.0));
}
i
}
#[test]
fn a_frame_round_trips_in_a_few_dozen_bytes() {
let frame = Frame::new(
183_456,
VizLevels {
low: 0.5126,
mid: 0.3,
high: 2.0,
},
);
let json = serde_json::to_string(&frame).unwrap();
assert_eq!(json, "[183456,513,300,1000]");
assert_eq!(serde_json::from_str::<Frame>(&json).unwrap(), frame);
assert!((frame.levels().low - 0.513).abs() < 1e-6);
}
#[test]
fn levels_between_frames_are_interpolated_behind_the_playhead() {
let mut i = interp(10_000, 10);
let delay = i.delay_ms();
assert!((34..=40).contains(&delay), "about two frames: {delay}");
let got = i.sample(10_050 + delay, true, Instant::now());
assert!((got.low - 0.025).abs() < 1e-4, "{got:?}");
assert!((got.high - 0.025 / 4.0).abs() < 1e-4);
}
#[test]
fn a_dropped_frame_is_bridged_by_the_delay() {
let mut i = Interp::default();
for n in [0u64, 1, 2, 4, 5] {
i.push(n * 20, lv(n as f32 / 10.0));
}
let got = i.sample(70 + i.delay_ms(), true, Instant::now());
assert!((got.low - 0.35).abs() < 1e-3, "{got:?}");
}
#[test]
fn a_dry_buffer_eases_to_rest_without_inventing_motion() {
let mut i = interp(0, 5);
let start = Instant::now();
let last = i.sample(80 + i.delay_ms(), true, start);
assert!(last.low > 0.0);
let mut prev = last.low;
for n in 1..=60u64 {
let got = i.sample(
80 + i.delay_ms() + n * 20,
true,
start + Duration::from_millis(n * 20),
);
assert!(got.low <= prev, "only ever falls");
prev = got.low;
}
assert!(prev < REST, "at rest: {prev}");
assert!(!i.moving());
}
#[test]
fn a_paused_device_eases_to_rest() {
let mut i = interp(0, 10);
let start = Instant::now();
i.sample(100 + i.delay_ms(), true, start);
let got = i.sample(100 + i.delay_ms(), false, start + Duration::from_secs(1));
assert!(got.low < REST);
}
#[test]
fn a_seek_clears_the_buffer() {
let mut i = interp(30_000, 10);
i.push(5_000, lv(0.9));
assert_eq!(i.frames.len(), 1, "only the frame after the seek");
i.push(5_020, lv(0.7));
let got = i.sample(5_010 + i.delay_ms(), true, Instant::now());
assert!((got.low - 0.8).abs() < 1e-4);
}
#[test]
fn a_frame_where_the_playhead_stood_replaces_the_last() {
let mut i = interp(0, 3);
i.push(40, lv(0.0));
assert_eq!(i.frames.len(), 3);
assert_eq!(i.frames.back().unwrap().1, lv(0.0));
}
#[test]
fn nothing_is_sent_with_no_watcher_and_it_stops_on_unsubscribe() {
let feed = Feed::new();
let waker = Waker::new().unwrap();
feed.publish_for_test(Frame(1, 1, 1, 1));
assert_eq!(feed.watchers(), 0);
let mut watch = feed.watch(&waker);
assert_eq!(watch.take(), None, "nothing new since it subscribed");
feed.publish_for_test(Frame(2, 2, 2, 2));
assert_eq!(watch.take(), Some(Frame(2, 2, 2, 2)));
assert_eq!(watch.take(), None, "sent once");
feed.publish_for_test(Frame(3, 3, 3, 3));
feed.publish_for_test(Frame(4, 4, 4, 4));
assert_eq!(
watch.take(),
Some(Frame(4, 4, 4, 4)),
"the newest, not a backlog"
);
drop(watch);
assert_eq!(feed.watchers(), 0);
}
#[test]
fn the_feed_reads_the_analyser_only_while_watched() {
let feed = Feed::new();
let viz = VizSnapshot::new();
feed.provide(Arc::clone(&viz), || 1_234);
let waker = Waker::new().unwrap();
let mut watch = feed.watch(&waker);
let deadline = Instant::now() + Duration::from_secs(5);
let frame = loop {
viz.write(Default::default());
if let Some(frame) = watch.take() {
break frame;
}
assert!(Instant::now() < deadline, "no frame reached the watcher");
std::thread::sleep(Duration::from_millis(5));
};
assert_eq!(frame.at_ms(), 1_234);
drop(watch);
drop(waker);
let reads = viz.reads();
for _ in 0..5 {
viz.write(Default::default());
std::thread::sleep(Duration::from_millis(5));
}
assert!(
viz.reads() <= reads + 1,
"the pump stopped reading once nobody watched"
);
}
}