1use std::collections::VecDeque;
19use std::sync::atomic::{AtomicU64, Ordering};
20use std::sync::{Arc, LazyLock, Weak};
21use std::time::{Duration, Instant};
22
23use parking_lot::{Condvar, Mutex};
24use serde::{Deserialize, Serialize};
25
26use crate::audio::viz::{VizLevels, VizSnapshot};
27use crate::remote::link::LinkCommand;
28use crate::remote::wire::Waker;
29use crate::signal::Wake;
30
31#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
35pub struct Frame(pub u64, pub u16, pub u16, pub u16);
36
37impl Frame {
38 pub fn new(at_ms: u64, levels: VizLevels) -> Self {
39 let q = |v: f32| (v.clamp(0.0, 1.0) * 1000.0).round() as u16;
40 Self(at_ms, q(levels.low), q(levels.mid), q(levels.high))
41 }
42
43 pub fn at_ms(&self) -> u64 {
44 self.0
45 }
46
47 pub fn levels(&self) -> VizLevels {
48 VizLevels {
49 low: f32::from(self.1) / 1000.0,
50 mid: f32::from(self.2) / 1000.0,
51 high: f32::from(self.3) / 1000.0,
52 }
53 }
54}
55
56pub struct Feed {
64 inner: Mutex<FeedState>,
65 changed: Condvar,
66 next_id: AtomicU64,
67}
68
69struct FeedState {
70 source: Option<Arc<Source>>,
71 watchers: Vec<(u64, Weak<Waker>)>,
72 latest: Option<(u64, Frame)>,
75 pumping: bool,
76}
77
78struct Source {
79 viz: Arc<VizSnapshot>,
80 position_ms: Box<dyn Fn() -> u64 + Send + Sync>,
81}
82
83static FEED: LazyLock<Arc<Feed>> = LazyLock::new(Feed::new);
84
85pub fn feed() -> &'static Arc<Feed> {
87 &FEED
88}
89
90impl Feed {
91 pub fn new() -> Arc<Self> {
92 Arc::new(Self {
93 inner: Mutex::new(FeedState {
94 source: None,
95 watchers: Vec::new(),
96 latest: None,
97 pumping: false,
98 }),
99 changed: Condvar::new(),
100 next_id: AtomicU64::new(1),
101 })
102 }
103
104 pub fn provide(
106 self: &Arc<Self>,
107 viz: Arc<VizSnapshot>,
108 position_ms: impl Fn() -> u64 + Send + Sync + 'static,
109 ) {
110 let mut inner = self.inner.lock();
111 inner.source = Some(Arc::new(Source {
112 viz,
113 position_ms: Box::new(position_ms),
114 }));
115 self.changed.notify_all();
116 }
117
118 pub fn watch(self: &Arc<Self>, waker: &Arc<Waker>) -> Watch {
121 let id = self.next_id.fetch_add(1, Ordering::Relaxed);
122 let mut inner = self.inner.lock();
123 inner.watchers.push((id, Arc::downgrade(waker)));
124 let sent = inner.latest.map_or(0, |(seq, _)| seq);
125 if !inner.pumping {
126 inner.pumping = true;
127 let feed = Arc::clone(self);
128 std::thread::Builder::new()
129 .name("koan-levels".into())
130 .spawn(move || feed.pump())
131 .expect("failed to spawn the levels thread");
132 }
133 self.changed.notify_all();
134 Watch {
135 feed: Arc::downgrade(self),
136 id,
137 sent,
138 }
139 }
140
141 pub fn watchers(&self) -> usize {
143 self.inner.lock().watchers.len()
144 }
145
146 fn pump(&self) {
147 loop {
148 let source = {
149 let mut inner = self.inner.lock();
150 loop {
151 inner.watchers.retain(|(_, w)| w.strong_count() > 0);
152 match &inner.source {
153 Some(source) if !inner.watchers.is_empty() => break Arc::clone(source),
154 _ => self.changed.wait(&mut inner),
155 }
156 }
157 };
158 let seen = source.viz.frames().generation();
160 source.viz.touch();
161 source.viz.frames().wait(seen);
162 let frame = Frame::new((source.position_ms)(), source.viz.levels());
163 let mut inner = self.inner.lock();
164 let seq = inner.latest.map_or(0, |(seq, _)| seq) + 1;
165 inner.latest = Some((seq, frame));
166 inner.watchers.retain(|(_, w)| match w.upgrade() {
167 Some(w) => {
168 w.wake();
169 true
170 }
171 None => false,
172 });
173 }
174 }
175
176 #[cfg(test)]
177 fn publish_for_test(&self, frame: Frame) {
178 let mut inner = self.inner.lock();
179 let seq = inner.latest.map_or(0, |(seq, _)| seq) + 1;
180 inner.latest = Some((seq, frame));
181 }
182}
183
184pub struct Watch {
186 feed: Weak<Feed>,
187 id: u64,
188 sent: u64,
189}
190
191impl Watch {
192 pub fn take(&mut self) -> Option<Frame> {
195 let feed = self.feed.upgrade()?;
196 let (seq, frame) = feed.inner.lock().latest?;
197 (seq > self.sent).then(|| {
198 self.sent = seq;
199 frame
200 })
201 }
202}
203
204impl Drop for Watch {
205 fn drop(&mut self) {
206 if let Some(feed) = self.feed.upgrade() {
207 feed.inner.lock().watchers.retain(|(id, _)| *id != self.id);
208 }
209 }
210}
211
212const DELAY_FRAMES: f32 = 2.0;
217const JUMP_MS: u64 = 1_000;
220const BACK_MS: u64 = 50;
222const EASE_HALF_LIFE: Duration = Duration::from_millis(80);
224const REST: f32 = 0.001;
226const MAX_FRAMES: usize = 240;
227
228#[derive(Debug)]
230pub struct Interp {
231 frames: VecDeque<(u64, VizLevels)>,
232 spacing_ms: f32,
235 shown: VizLevels,
236 sampled: Option<Instant>,
237}
238
239impl Default for Interp {
240 fn default() -> Self {
241 Self {
242 frames: VecDeque::new(),
243 spacing_ms: 1000.0 / 60.0,
244 shown: VizLevels::default(),
245 sampled: None,
246 }
247 }
248}
249
250impl Interp {
251 pub fn push(&mut self, at_ms: u64, levels: VizLevels) {
255 if let Some(&(newest, _)) = self.frames.back() {
256 if at_ms + BACK_MS < newest || at_ms > newest + JUMP_MS {
259 self.frames.clear();
260 } else if at_ms <= newest {
261 if let Some(back) = self.frames.back_mut() {
262 *back = (newest, levels);
263 }
264 return;
265 } else {
266 let gap = (at_ms - newest) as f32;
267 if gap < 200.0 {
268 self.spacing_ms = self.spacing_ms * 0.8 + gap * 0.2;
269 }
270 }
271 }
272 self.frames.push_back((at_ms, levels));
273 while self.frames.len() > MAX_FRAMES {
274 self.frames.pop_front();
275 }
276 }
277
278 pub fn clear(&mut self) {
280 self.frames.clear();
281 }
282
283 pub fn delay_ms(&self) -> u64 {
285 (self.spacing_ms * DELAY_FRAMES).round() as u64
286 }
287
288 pub fn sample(&mut self, playhead_ms: u64, playing: bool, now: Instant) -> VizLevels {
292 let elapsed = self
293 .sampled
294 .map_or(Duration::ZERO, |at| now.saturating_duration_since(at));
295 self.sampled = Some(now);
296 let at = playhead_ms.saturating_sub(self.delay_ms());
297 match self.between(at).filter(|_| playing) {
298 Some(levels) => self.shown = levels,
299 None => {
300 let keep = 0.5f32.powf(elapsed.as_secs_f32() / EASE_HALF_LIFE.as_secs_f32());
301 self.shown = VizLevels {
302 low: self.shown.low * keep,
303 mid: self.shown.mid * keep,
304 high: self.shown.high * keep,
305 };
306 let stale = self
308 .frames
309 .back()
310 .is_some_and(|&(newest, _)| at > newest + JUMP_MS || at + JUMP_MS < newest);
311 if stale || !playing {
312 self.frames.clear();
313 }
314 }
315 }
316 self.shown
317 }
318
319 fn between(&mut self, at: u64) -> Option<VizLevels> {
322 let after = self.frames.iter().position(|&(t, _)| t >= at)?;
323 if after == 0 {
324 return (self.frames[0].0 == at).then_some(self.frames[0].1);
325 }
326 self.frames.drain(..after - 1);
327 let (t0, a) = self.frames[0];
328 let (t1, b) = self.frames[1];
329 let f = (at - t0) as f32 / (t1 - t0) as f32;
330 let lerp = |x: f32, y: f32| x + (y - x) * f;
331 Some(VizLevels {
332 low: lerp(a.low, b.low),
333 mid: lerp(a.mid, b.mid),
334 high: lerp(a.high, b.high),
335 })
336 }
337
338 pub fn moving(&self) -> bool {
341 !self.frames.is_empty()
342 || self.shown.low > REST
343 || self.shown.mid > REST
344 || self.shown.high > REST
345 }
346}
347
348const STALL: Duration = Duration::from_secs(3);
351
352pub struct Remote {
357 inner: Mutex<RemoteState>,
358 ticks: Wake,
361 tick: Condvar,
362 send: Box<SendWatch>,
366}
367
368type SendWatch = dyn Fn(&str, LinkCommand) -> bool + Send + Sync;
370
371struct RemoteState {
372 views: Vec<Viewed>,
373 drawing: Option<String>,
375 interp: Interp,
376 heard: Option<Instant>,
378 interval: Duration,
379 ticking: bool,
380}
381
382struct Viewed {
384 target: String,
385 count: usize,
386 asked: Option<u64>,
391}
392
393static REMOTE: LazyLock<Arc<Remote>> =
394 LazyLock::new(|| Remote::new(crate::remote::devices::send_live));
395
396pub fn remote() -> &'static Arc<Remote> {
397 &REMOTE
398}
399
400impl Remote {
401 pub fn new(send: impl Fn(&str, LinkCommand) -> bool + Send + Sync + 'static) -> Arc<Self> {
402 Arc::new(Self {
403 inner: Mutex::new(RemoteState {
404 views: Vec::new(),
405 drawing: None,
406 interp: Interp::default(),
407 heard: None,
408 interval: Duration::from_micros(1_000_000 / 60),
409 ticking: false,
410 }),
411 ticks: Wake::new(),
412 tick: Condvar::new(),
413 send: Box::new(send),
414 })
415 }
416
417 pub fn view(self: &Arc<Self>, target: String) -> View {
419 let mut inner = self.inner.lock();
420 match inner.views.iter_mut().find(|v| v.target == target) {
421 Some(v) => v.count += 1,
422 None => inner.views.push(Viewed {
423 target: target.clone(),
424 count: 1,
425 asked: None,
426 }),
427 }
428 if inner.drawing.as_deref() != Some(&target) {
429 inner.drawing = Some(target.clone());
430 inner.interp = Interp::default();
431 inner.heard = None;
432 }
433 if !inner.ticking {
434 inner.ticking = true;
435 let remote = Arc::clone(self);
436 std::thread::Builder::new()
437 .name("koan-levels-tick".into())
438 .spawn(move || remote.run_ticks())
439 .expect("failed to spawn the levels ticker");
440 }
441 self.tick.notify_all();
442 drop(inner);
443 self.ask(false);
444 View {
445 remote: Arc::downgrade(self),
446 target,
447 }
448 }
449
450 pub fn ask(&self, again: bool) {
455 let version = crate::remote::devices::version();
456 let due: Vec<String> = {
457 let inner = self.inner.lock();
458 inner
459 .views
460 .iter()
461 .filter(|v| again || v.asked != Some(version))
462 .map(|v| v.target.clone())
463 .collect()
464 };
465 for target in due {
466 let sent = (self.send)(&target, LinkCommand::WatchLevels { on: true });
467 if !sent {
468 log::debug!("levels: no way to {target} now; asking when one opens");
469 }
470 if let Some(v) = self
471 .inner
472 .lock()
473 .views
474 .iter_mut()
475 .find(|v| v.target == target)
476 {
477 v.asked = sent.then_some(version);
478 }
479 }
480 }
481
482 pub fn received(&self, from: &str, frame: Frame) {
484 let mut inner = self.inner.lock();
485 if inner.drawing.as_deref() != Some(from) || !inner.views.iter().any(|v| v.target == from) {
486 return;
487 }
488 inner.interp.push(frame.at_ms(), frame.levels());
489 inner.heard = Some(Instant::now());
490 self.tick.notify_all();
491 }
492
493 pub fn sample(&self, playhead_ms: u64, playing: bool) -> VizLevels {
495 self.inner
496 .lock()
497 .interp
498 .sample(playhead_ms, playing, Instant::now())
499 }
500
501 pub fn set_fps(&self, fps: u8) {
503 self.inner.lock().interval = Duration::from_micros(1_000_000 / fps.clamp(1, 240) as u64);
504 }
505
506 pub fn ticks(&self) -> &Wake {
508 &self.ticks
509 }
510
511 fn release(&self, target: &str) {
512 let mut inner = self.inner.lock();
513 let Some(at) = inner.views.iter().position(|v| v.target == target) else {
514 return;
515 };
516 inner.views[at].count -= 1;
517 if inner.views[at].count > 0 {
518 return;
519 }
520 inner.views.remove(at);
521 if inner.drawing.as_deref() == Some(target) {
522 inner.drawing = None;
523 inner.interp = Interp::default();
524 inner.heard = None;
525 }
526 drop(inner);
527 (self.send)(target, LinkCommand::WatchLevels { on: false });
528 }
529
530 fn run_ticks(&self) {
531 loop {
532 let interval = {
533 let mut inner = self.inner.lock();
534 loop {
535 if inner.views.is_empty() {
536 self.tick.wait(&mut inner);
537 } else if inner.interp.moving() {
538 break inner.interval;
539 } else if self.tick.wait_for(&mut inner, STALL).timed_out()
540 && inner.heard.is_none_or(|at| at.elapsed() >= STALL)
541 && !inner.views.is_empty()
542 {
543 drop(inner);
544 let playing = crate::remote::devices::target_playhead()
545 .is_some_and(|(_, playing)| playing);
546 if playing {
547 self.ask(true);
548 }
549 inner = self.inner.lock();
550 }
551 }
552 };
553 std::thread::sleep(interval);
554 self.ticks.bump();
555 }
556 }
557
558 #[cfg(test)]
559 fn viewers(&self, target: &str) -> usize {
560 self.inner
561 .lock()
562 .views
563 .iter()
564 .find(|v| v.target == target)
565 .map_or(0, |v| v.count)
566 }
567}
568
569pub struct View {
571 remote: Weak<Remote>,
572 target: String,
573}
574
575impl View {
576 pub fn target(&self) -> &str {
577 &self.target
578 }
579}
580
581impl Drop for View {
582 fn drop(&mut self) {
583 if let Some(remote) = self.remote.upgrade() {
584 remote.release(&self.target);
585 }
586 }
587}
588
589#[cfg(test)]
590mod tests {
591 use super::*;
592
593 type Sent = Arc<Mutex<Vec<(String, bool)>>>;
595
596 fn recording(through: bool) -> (Arc<Remote>, Sent, Arc<std::sync::atomic::AtomicBool>) {
597 let sent = Arc::new(Mutex::new(Vec::new()));
598 let open = Arc::new(std::sync::atomic::AtomicBool::new(through));
599 let (log, gate) = (Arc::clone(&sent), Arc::clone(&open));
600 let remote = Remote::new(move |to, cmd| {
601 if let LinkCommand::WatchLevels { on } = cmd {
602 log.lock().push((to.to_string(), on));
603 }
604 gate.load(Ordering::Relaxed)
605 });
606 (remote, sent, open)
607 }
608
609 #[test]
610 fn switching_devices_tells_the_old_one_to_stop() {
611 let (remote, sent, _) = recording(true);
612 let x = remote.view("x".into());
613 let y = remote.view("y".into());
615 drop(x);
616 assert_eq!(remote.viewers("x"), 0);
617 assert!(
618 sent.lock().contains(&("x".into(), false)),
619 "{:?}",
620 sent.lock()
621 );
622 assert!(!sent.lock().contains(&("y".into(), false)));
623 drop(y);
624 assert!(sent.lock().contains(&("y".into(), false)));
625 }
626
627 #[test]
628 fn a_device_stays_watched_while_any_view_of_it_is_held() {
629 let (remote, sent, _) = recording(true);
630 let a = remote.view("x".into());
631 let b = remote.view("x".into());
632 drop(a);
633 assert!(!sent.lock().contains(&("x".into(), false)));
634 drop(b);
635 assert!(sent.lock().contains(&("x".into(), false)));
636 }
637
638 #[test]
639 fn an_ask_that_found_no_way_through_is_made_again() {
640 let (remote, sent, open) = recording(false);
641 let _view = remote.view("x".into());
642 let asks = |sent: &Mutex<Vec<(String, bool)>>| {
643 sent.lock().iter().filter(|(t, on)| t == "x" && *on).count()
644 };
645 assert_eq!(asks(&sent), 1);
646 open.store(true, Ordering::Relaxed);
648 remote.ask(false);
649 assert_eq!(asks(&sent), 2, "asked again, having not got through");
650 remote.ask(true);
651 assert_eq!(asks(&sent), 3, "and on a renewal");
652 }
653
654 #[test]
655 fn frames_from_a_device_not_viewed_are_ignored() {
656 let (remote, _, _) = recording(true);
657 let _view = remote.view("x".into());
658 remote.received("y", Frame(100, 500, 500, 500));
659 assert!(!remote.inner.lock().interp.moving());
660 remote.received("x", Frame(100, 500, 500, 500));
661 assert!(remote.inner.lock().interp.moving());
662 }
663
664 #[test]
665 fn a_short_skip_back_clears_the_buffer() {
666 let mut i = interp(400, 10);
667 i.push(150, lv(0.9));
668 assert_eq!(i.frames.len(), 1);
669 }
670
671 fn lv(v: f32) -> VizLevels {
672 VizLevels {
673 low: v,
674 mid: v / 2.0,
675 high: v / 4.0,
676 }
677 }
678
679 fn interp(from: u64, frames: u64) -> Interp {
681 let mut i = Interp::default();
682 for n in 0..frames {
683 i.push(from + n * 20, lv(n as f32 / 100.0));
684 }
685 i
686 }
687
688 #[test]
689 fn a_frame_round_trips_in_a_few_dozen_bytes() {
690 let frame = Frame::new(
691 183_456,
692 VizLevels {
693 low: 0.5126,
694 mid: 0.3,
695 high: 2.0,
696 },
697 );
698 let json = serde_json::to_string(&frame).unwrap();
699 assert_eq!(json, "[183456,513,300,1000]");
700 assert_eq!(serde_json::from_str::<Frame>(&json).unwrap(), frame);
701 assert!((frame.levels().low - 0.513).abs() < 1e-6);
702 }
703
704 #[test]
705 fn levels_between_frames_are_interpolated_behind_the_playhead() {
706 let mut i = interp(10_000, 10);
707 let delay = i.delay_ms();
708 assert!((34..=40).contains(&delay), "about two frames: {delay}");
709 let got = i.sample(10_050 + delay, true, Instant::now());
711 assert!((got.low - 0.025).abs() < 1e-4, "{got:?}");
712 assert!((got.high - 0.025 / 4.0).abs() < 1e-4);
713 }
714
715 #[test]
716 fn a_dropped_frame_is_bridged_by_the_delay() {
717 let mut i = Interp::default();
718 for n in [0u64, 1, 2, 4, 5] {
719 i.push(n * 20, lv(n as f32 / 10.0));
720 }
721 let got = i.sample(70 + i.delay_ms(), true, Instant::now());
723 assert!((got.low - 0.35).abs() < 1e-3, "{got:?}");
724 }
725
726 #[test]
727 fn a_dry_buffer_eases_to_rest_without_inventing_motion() {
728 let mut i = interp(0, 5);
729 let start = Instant::now();
730 let last = i.sample(80 + i.delay_ms(), true, start);
731 assert!(last.low > 0.0);
732 let mut prev = last.low;
734 for n in 1..=60u64 {
735 let got = i.sample(
736 80 + i.delay_ms() + n * 20,
737 true,
738 start + Duration::from_millis(n * 20),
739 );
740 assert!(got.low <= prev, "only ever falls");
741 prev = got.low;
742 }
743 assert!(prev < REST, "at rest: {prev}");
744 assert!(!i.moving());
745 }
746
747 #[test]
748 fn a_paused_device_eases_to_rest() {
749 let mut i = interp(0, 10);
750 let start = Instant::now();
751 i.sample(100 + i.delay_ms(), true, start);
752 let got = i.sample(100 + i.delay_ms(), false, start + Duration::from_secs(1));
753 assert!(got.low < REST);
754 }
755
756 #[test]
757 fn a_seek_clears_the_buffer() {
758 let mut i = interp(30_000, 10);
759 i.push(5_000, lv(0.9));
760 assert_eq!(i.frames.len(), 1, "only the frame after the seek");
761 i.push(5_020, lv(0.7));
762 let got = i.sample(5_010 + i.delay_ms(), true, Instant::now());
763 assert!((got.low - 0.8).abs() < 1e-4);
764 }
765
766 #[test]
767 fn a_frame_where_the_playhead_stood_replaces_the_last() {
768 let mut i = interp(0, 3);
769 i.push(40, lv(0.0));
770 assert_eq!(i.frames.len(), 3);
771 assert_eq!(i.frames.back().unwrap().1, lv(0.0));
772 }
773
774 #[test]
775 fn nothing_is_sent_with_no_watcher_and_it_stops_on_unsubscribe() {
776 let feed = Feed::new();
777 let waker = Waker::new().unwrap();
778 feed.publish_for_test(Frame(1, 1, 1, 1));
779 assert_eq!(feed.watchers(), 0);
780
781 let mut watch = feed.watch(&waker);
782 assert_eq!(watch.take(), None, "nothing new since it subscribed");
783 feed.publish_for_test(Frame(2, 2, 2, 2));
784 assert_eq!(watch.take(), Some(Frame(2, 2, 2, 2)));
785 assert_eq!(watch.take(), None, "sent once");
786 feed.publish_for_test(Frame(3, 3, 3, 3));
787 feed.publish_for_test(Frame(4, 4, 4, 4));
788 assert_eq!(
789 watch.take(),
790 Some(Frame(4, 4, 4, 4)),
791 "the newest, not a backlog"
792 );
793
794 drop(watch);
795 assert_eq!(feed.watchers(), 0);
796 }
797
798 #[test]
799 fn the_feed_reads_the_analyser_only_while_watched() {
800 let feed = Feed::new();
801 let viz = VizSnapshot::new();
802 feed.provide(Arc::clone(&viz), || 1_234);
803 let waker = Waker::new().unwrap();
804 let mut watch = feed.watch(&waker);
805
806 let deadline = Instant::now() + Duration::from_secs(5);
807 let frame = loop {
808 viz.write(Default::default());
809 if let Some(frame) = watch.take() {
810 break frame;
811 }
812 assert!(Instant::now() < deadline, "no frame reached the watcher");
813 std::thread::sleep(Duration::from_millis(5));
814 };
815 assert_eq!(frame.at_ms(), 1_234);
816
817 drop(watch);
819 drop(waker);
820 let reads = viz.reads();
821 for _ in 0..5 {
822 viz.write(Default::default());
823 std::thread::sleep(Duration::from_millis(5));
824 }
825 assert!(
826 viz.reads() <= reads + 1,
827 "the pump stopped reading once nobody watched"
828 );
829 }
830}