From 96d323c32de45da159d76f033c9179365bcc81c7 Mon Sep 17 00:00:00 2001 From: Mikei386 <44135113+Mikei386@users.noreply.github.com> Date: Wed, 5 Aug 2026 20:53:58 +0200 Subject: [PATCH] Revert "Decouple realtime snapshots from transport cadence" This reverts commit feec2bac56c741122386345ea46ce85034e4d1c1. --- src/audio.rs | 21 +++++++++--- src/goniometer.rs | 27 +++++++++++++-- src/routes.rs | 86 ++++++----------------------------------------- 3 files changed, 52 insertions(+), 82 deletions(-) diff --git a/src/audio.rs b/src/audio.rs index ceb26f7..9cc1a85 100644 --- a/src/audio.rs +++ b/src/audio.rs @@ -14,7 +14,7 @@ use tracing::warn; #[cfg(target_os = "linux")] use crate::correlation::CorrelationMeter; #[cfg(target_os = "linux")] -use crate::goniometer::selected_sample_indices; +use crate::goniometer::{selected_sample_indices, GoniometerClock}; #[cfg(target_os = "linux")] use crate::model::{RtaFrame, SpectroFrame, WaveEnvFrame}; #[cfg(target_os = "linux")] @@ -40,8 +40,12 @@ use crate::{ model::{InputSource, MeterFrame, PhoenixRtaConfig}, }; -#[cfg(not(target_os = "linux"))] -const PLACEHOLDER_UPDATES_PER_SECOND: u32 = 120; +// Keep transport sampling safely above the 60 Hz display cadence. ALSA only +// exposes complete capture periods, so a 60 Hz target with 128-frame periods +// alternates between roughly 16.0 and 18.7 ms at 48 kHz. A 120 Hz snapshot +// clock guarantees at least one fresh native state per display frame without +// changing any detector or ballistic calculation. +const METRICS_TARGET_UPDATES_PER_SECOND: u32 = 120; #[derive(Clone, Copy, Debug, Default)] struct TransportPeaks { @@ -415,7 +419,7 @@ pub fn spawn_audio_capture_worker(deps: AudioWorkerDeps) { pub fn spawn_audio_capture_worker(deps: AudioWorkerDeps) { tokio::spawn(async move { warn!("Phoenix ALSA capture is only available on Linux; emitting placeholder frames on this host"); - let tick_ns = 1_000_000_000u64 / u64::from(PLACEHOLDER_UPDATES_PER_SECOND); + let tick_ns = 1_000_000_000u64 / u64::from(METRICS_TARGET_UPDATES_PER_SECOND); let mut ticker = tokio::time::interval(Duration::from_nanos(tick_ns)); deps.actual_sample_rate .store(deps.config.sample_rate as u64, Ordering::SeqCst); @@ -594,6 +598,7 @@ struct PpmState { lr_delay_y1_r: f32, true_peak_l: TruePeakDetector, true_peak_r: TruePeakDetector, + transport_clock: GoniometerClock, transport_peaks: TransportPeaks, correlation: CorrelationMeter, phase_wheel: PhaseWheelAnalyzer, @@ -626,6 +631,7 @@ impl Default for PpmState { lr_delay_y1_r: 0.0, true_peak_l: TruePeakDetector::default(), true_peak_r: TruePeakDetector::default(), + transport_clock: GoniometerClock::default(), transport_peaks: TransportPeaks::default(), correlation: CorrelationMeter::new(48_000, 1.0, -75.0, 0), phase_wheel: PhaseWheelAnalyzer::new(48_000), @@ -1251,6 +1257,13 @@ fn process_audio_block( None }; + if !ppm_state + .transport_clock + .advance(frames, sample_rate, METRICS_TARGET_UPDATES_PER_SECOND) + { + return (None, spectro_frame); + } + let rta = ppm_state.rta_state.as_ref().and_then(|state| match state { RtaEngineState::Iir(bank) => Some(build_iir_rta_frame(bank, sample_rate, rta_config)), RtaEngineState::Fft(fft) if fft.ring_fill >= fft.fft_size => { diff --git a/src/goniometer.rs b/src/goniometer.rs index efbea66..1d3543c 100644 --- a/src/goniometer.rs +++ b/src/goniometer.rs @@ -1,12 +1,10 @@ //! Timing and sample-selection helpers for the realtime goniometer stream. -#[cfg(test)] #[derive(Clone, Debug, Default)] pub struct GoniometerClock { phase: u64, } -#[cfg(test)] impl GoniometerClock { /// Returns true at the first capture boundary after the next display tick. /// The fractional phase is retained, so the average rate is independent @@ -60,6 +58,31 @@ mod tests { } } + #[test] + fn double_rate_transport_never_leaves_a_sixty_hz_display_interval_empty() { + let sample_rate = 48_000u32; + let period = 128usize; + let mut clock = GoniometerClock::default(); + let mut processed = 0usize; + let mut previous_emission = None; + let mut largest_gap = 0usize; + + while processed < sample_rate as usize * 10 { + processed += period; + if clock.advance(period, sample_rate, 120) { + if let Some(previous) = previous_emission { + largest_gap = largest_gap.max(processed - previous); + } + previous_emission = Some(processed); + } + } + + assert!( + largest_gap <= sample_rate as usize / 60, + "largest transport gap was {largest_gap} samples" + ); + } + #[test] fn selection_never_upsamples_and_keeps_endpoints() { assert_eq!(selected_sample_indices(3, 1024), vec![0, 1, 2]); diff --git a/src/routes.rs b/src/routes.rs index 9fb5c9d..4de525e 100644 --- a/src/routes.rs +++ b/src/routes.rs @@ -14,7 +14,7 @@ use std::{ process::Stdio, sync::atomic::{AtomicU64, Ordering}, sync::Arc, - time::{Duration, SystemTime, UNIX_EPOCH}, + time::{SystemTime, UNIX_EPOCH}, }; use tokio::{fs, io::AsyncWriteExt, process::Command}; use tracing::warn; @@ -27,8 +27,6 @@ use crate::{ const ONLINE_UPDATE_ZIP_URL: &str = "https://webshare.casaderoll.de/share/Phoenix.zip"; const DEVICE_BACKUP_KIND: &str = "phoenix-device-backup"; const DEVICE_BACKUP_SCHEMA_VERSION: u32 = 1; -const REALTIME_WS_UPDATES_PER_SECOND: u64 = 120; -const MAX_MERGED_XY_POINTS: usize = 2048; static JSON_WRITE_TOKEN: AtomicU64 = AtomicU64::new(1); fn stable_script_command(program: &str) -> Command { @@ -1198,11 +1196,7 @@ fn normalize_frontend_layout_id(raw: &str) -> String { async fn metrics_ws_inner(mut socket: axum::extract::ws::WebSocket, state: AppState) { let mut rx = state.subscribe_metrics(); - let tick_ns = 1_000_000_000u64 / REALTIME_WS_UPDATES_PER_SECOND; - let mut ticker = tokio::time::interval(Duration::from_nanos(tick_ns)); - ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { - ticker.tick().await; let Some(mut latest) = recv_latest_meter_frame(&mut rx).await else { break; }; @@ -1242,11 +1236,7 @@ fn strip_visual_payloads(frame: &mut MeterFrame) { async fn visuals_ws_inner(mut socket: axum::extract::ws::WebSocket, state: AppState) { let mut rx = state.subscribe_metrics(); - let tick_ns = 1_000_000_000u64 / REALTIME_WS_UPDATES_PER_SECOND; - let mut ticker = tokio::time::interval(Duration::from_nanos(tick_ns)); - ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { - ticker.tick().await; let Some(mut latest) = recv_latest_meter_frame(&mut rx).await else { return; }; @@ -1426,14 +1416,19 @@ fn drain_visual_meter_frames( latest: &mut Arc, ) -> (Option, Option<(Vec, Vec)>) { let mut combined_wave_env = latest.wave_env.clone(); - let mut combined_xy = None; - append_xy_samples(&mut combined_xy, &latest.xy_l, &latest.xy_r); + let mut latest_xy = if latest.xy_l.is_empty() || latest.xy_r.is_empty() { + None + } else { + Some((latest.xy_l.clone(), latest.xy_r.clone())) + }; loop { match rx.try_recv() { Ok(newer) => { merge_wave_env(&mut combined_wave_env, newer.wave_env.clone()); - append_xy_samples(&mut combined_xy, &newer.xy_l, &newer.xy_r); + if !newer.xy_l.is_empty() && !newer.xy_r.is_empty() { + latest_xy = Some((newer.xy_l.clone(), newer.xy_r.clone())); + } *latest = newer; } Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break, @@ -1444,31 +1439,7 @@ fn drain_visual_meter_frames( } } - (combined_wave_env, combined_xy) -} - -fn append_xy_samples( - target: &mut Option<(Vec, Vec)>, - incoming_l: &[f32], - incoming_r: &[f32], -) { - let count = incoming_l.len().min(incoming_r.len()); - if count == 0 { - return; - } - let (left, right) = target.get_or_insert_with(|| { - ( - Vec::with_capacity(count.min(MAX_MERGED_XY_POINTS)), - Vec::with_capacity(count.min(MAX_MERGED_XY_POINTS)), - ) - }); - left.extend_from_slice(&incoming_l[..count]); - right.extend_from_slice(&incoming_r[..count]); - if left.len() > MAX_MERGED_XY_POINTS { - let excess = left.len() - MAX_MERGED_XY_POINTS; - left.drain(..excess); - right.drain(..excess); - } + (combined_wave_env, latest_xy) } fn merge_wave_env(target: &mut Option, incoming: Option) { @@ -1718,43 +1689,6 @@ mod tests { assert_eq!(tp_r, -3.0); } - #[test] - fn visual_drain_keeps_xy_samples_in_chronological_order() { - let (tx, mut rx) = tokio::sync::broadcast::channel(8); - for (seq, left, right) in [ - (1, vec![1.0, 2.0], vec![-1.0, -2.0]), - (2, vec![3.0, 4.0], vec![-3.0, -4.0]), - (3, vec![5.0], vec![-5.0]), - ] { - let mut frame = meter_frame(); - frame.seq = seq; - frame.xy_l = left; - frame.xy_r = right; - tx.send(Arc::new(frame)).unwrap(); - } - - let mut latest = rx.try_recv().unwrap(); - let (_, xy) = drain_visual_meter_frames(&mut rx, &mut latest); - let (left, right) = xy.unwrap(); - assert_eq!(latest.seq, 3); - assert_eq!(left, vec![1.0, 2.0, 3.0, 4.0, 5.0]); - assert_eq!(right, vec![-1.0, -2.0, -3.0, -4.0, -5.0]); - } - - #[test] - fn merged_xy_keeps_the_newest_bounded_window() { - let mut xy = None; - let first: Vec = (0..1500).map(|value| value as f32).collect(); - let second: Vec = (1500..2500).map(|value| value as f32).collect(); - append_xy_samples(&mut xy, &first, &first); - append_xy_samples(&mut xy, &second, &second); - let (left, right) = xy.unwrap(); - assert_eq!(left.len(), MAX_MERGED_XY_POINTS); - assert_eq!(right.len(), MAX_MERGED_XY_POINTS); - assert_eq!(left[0], 452.0); - assert_eq!(left[MAX_MERGED_XY_POINTS - 1], 2499.0); - } - #[test] fn wave_envelopes_merge_without_losing_columns() { let mut target = Some(WaveEnvFrame {