diff --git a/TODO.md b/TODO.md index 83f84e3..465b46f 100644 --- a/TODO.md +++ b/TODO.md @@ -28,12 +28,13 @@ Diese vorhandenen Funktionen sind nicht automatisch messtechnisch korrekt. Die f - [x] **1. WebSocket auf "latest value wins" umstellen** - **Soll:** Ein langsamer Client erhält immer den neuesten Messzustand; alte Zustände werden verworfen. - - **Ist:** Der interne Kanal ist auf 32 Frames begrenzt; der WebSocket sendet im 16-ms-Takt und leert vor jedem Versand bis zum neuesten Zustand. Spektrogrammdaten laufen getrennt. - - **Abnahme:** Lokaler Laufzeittest liefert 68 aktuelle Messframes in 1,1 s; alte Zustände werden beim Leeren verworfen. + - **Ist:** Der DSP läuft weiterhin samplekontinuierlich, erzeugt skalare Transport-Snapshots aber periodengrößenunabhängig nur noch mit im Mittel 60 Hz. Der WebSocket sendet diese ohne zusätzlichen 16-ms-Ticker sofort weiter und leert bei Rückstau bis zum neuesten Zustand. + - **Peak-Schutz:** Sample- und True-Peak-Maxima werden über alle Capture-Blöcke bis zum nächsten Snapshot gesammelt. Muss der WebSocket mehrere Snapshots zusammenfassen, bleiben deren höchste True-Peak-Werte ebenfalls erhalten. + - **Abnahme:** Automatische Tests prüfen die periodengrößenunabhängige 60-Hz-Taktung, Peak-Erhalt und Latest-State-Semantik. - [x] **2. Mess-, Visualisierungs- und Konfigurationsdaten trennen** - **Soll:** DSP läuft samplegenau; übertragen wird nur so häufig und so umfangreich wie für die jeweilige Anzeige nötig. - - **Ist:** Skalare Messzustände laufen mit etwa 60 JSON-Paketen/s. Spektrogramm sowie Goniometer/Waveform besitzen getrennte, kompakte Binär-WebSockets. Roh-Waveformsamples werden nicht mehr in jedem Capture-Frame vervielfacht; Waveform-Hüllkurven werden beim serverseitigen Leeren lückenlos zusammengeführt. Das doppelte RTA-Bandfeld wurde vollständig entfernt; Konfiguration wird nur separat bei Änderungen synchronisiert. + - **Ist:** Skalare Messzustände laufen mit etwa 60 JSON-Paketen/s. Das Spektrogramm besitzt zusätzlich zum eigenen Binär-WebSocket nun auch einen vollständig getrennten Backend-Kanal und behält dadurch seinen FFT-Takt unabhängig von den Mess-Snapshots. Goniometer/Waveform verwenden einen kompakten Binär-WebSocket. Roh-Waveformsamples werden nicht mehr in jedem Capture-Frame vervielfacht; Waveform-Hüllkurven werden bereits zwischen zwei 60-Hz-Snapshots lückenlos gesammelt. Das doppelte RTA-Bandfeld wurde vollständig entfernt; Konfiguration wird nur separat bei Änderungen synchronisiert. - **Interne Last:** Der Verteiler reicht Messframes als gemeinsam genutzte Referenz weiter. Beim Verwerfen alter Frames werden deshalb keine kompletten Spektrogramm-, XY- und Waveformvektoren mehr kopiert. - **Geprüft:** Protokolltests prüfen Header, Nutzdaten und beschädigte Paketlängen. Ein Servertest stellt sicher, dass große Visualisierungsfelder nicht wieder im JSON-Messstrom landen. - **Abnahme:** Datenwege sind softwareseitig getrennt; die Lastmessung auf der Zielhardware bleibt Bestandteil der End-to-End-Abnahme unter Punkt 17. diff --git a/src/audio.rs b/src/audio.rs index 2aa7334..daf214b 100644 --- a/src/audio.rs +++ b/src/audio.rs @@ -37,9 +37,31 @@ use crate::{ }; #[cfg(target_os = "linux")] -// 128-sample periods at 48 kHz -> 62.5 visual updates/s. The DSP history is -// still updated sample-by-sample; only transport snapshots are rate-limited. -const XY_TARGET_UPDATES_PER_SECOND: u64 = 60; +const METRICS_TARGET_UPDATES_PER_SECOND: u32 = 60; + +#[derive(Clone, Copy, Debug, Default)] +struct TransportPeaks { + sample_l: f32, + sample_r: f32, + true_l: f32, + true_r: f32, +} + +impl TransportPeaks { + fn observe(&mut self, sample_l: f32, sample_r: f32, true_l: f32, true_r: f32) { + self.sample_l = self.sample_l.max(sample_l.abs()); + self.sample_r = self.sample_r.max(sample_r.abs()); + self.true_l = self.true_l.max(true_l.abs()); + self.true_r = self.true_r.max(true_r.abs()); + } + + fn take(&mut self) -> (f32, f32) { + let left = self.sample_l.max(self.true_l); + let right = self.sample_r.max(self.true_r); + *self = Self::default(); + (left, right) + } +} const RTA_PEAK_FLOOR_DB: f32 = -150.0; #[cfg(target_os = "linux")] const WAVE_ENV_COLUMNS_PER_SEC: f32 = 9600.0; @@ -99,6 +121,8 @@ pub struct AudioWorkerDeps { pub actual_sample_rate: Arc, pub metrics_tx: tokio::sync::broadcast::Sender>, #[cfg_attr(not(target_os = "linux"), allow(dead_code))] + pub spectro_tx: tokio::sync::broadcast::Sender>, + #[cfg_attr(not(target_os = "linux"), allow(dead_code))] pub spectro_subscribers: Arc, #[cfg_attr(not(target_os = "linux"), allow(dead_code))] pub restart_token: Arc, @@ -519,7 +543,7 @@ fn capture_until_restart(deps: AudioWorkerDeps, generation: u64) -> anyhow::Resu &buffer[..frames * 2], actual_sample_rate, ); - let frame = build_meter_frame( + let (metrics_frame, spectro_frame) = process_audio_block( &buffer[..frames * 2], actual_sample_rate, input, @@ -529,7 +553,12 @@ fn capture_until_restart(deps: AudioWorkerDeps, generation: u64) -> anyhow::Resu &rta_config, deps.spectro_subscribers.load(Ordering::Relaxed) > 0, ); - let _ = deps.metrics_tx.send(Arc::new(frame)); + if let Some(frame) = metrics_frame { + let _ = deps.metrics_tx.send(Arc::new(frame)); + } + if let Some(frame) = spectro_frame { + let _ = deps.spectro_tx.send(Arc::new(frame)); + } } Err(err) => match pcm.state() { PcmState::XRun | PcmState::Suspended => { @@ -550,8 +579,6 @@ struct PpmState { ebu_ppm: PpmDetector, vu_meter: VuMeter, rms_window: MovingAverageWindow, - last_rta: Option, - last_spectro: Option, rta_signature: String, rta_state: Option, spectro_signature: String, @@ -569,10 +596,11 @@ struct PpmState { lr_delay_y1_r: f32, true_peak_l: TruePeakDetector, true_peak_r: TruePeakDetector, + transport_clock: GoniometerClock, + transport_peaks: TransportPeaks, correlation: CorrelationMeter, xy_pending_l: Vec, xy_pending_r: Vec, - xy_clock: GoniometerClock, } #[cfg(target_os = "linux")] @@ -583,8 +611,6 @@ impl Default for PpmState { ebu_ppm: PpmDetector::new(48_000, PpmStandard::EbuTypeIib), vu_meter: VuMeter::new(48_000), rms_window: create_moving_average_window(48_000, VU_WINDOW_MS), - last_rta: None, - last_spectro: None, rta_signature: String::new(), rta_state: None, spectro_signature: String::new(), @@ -602,10 +628,11 @@ 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, 0), xy_pending_l: Vec::with_capacity(1024), xy_pending_r: Vec::with_capacity(1024), - xy_clock: GoniometerClock::default(), } } } @@ -1161,7 +1188,7 @@ fn take_goniometer_samples(state: &mut PpmState, target_points: usize) -> (Vec MeterFrame { +) -> (Option, Option) { let mut rms_power_l = 0.0f32; let mut rms_power_r = 0.0f32; - let mut peak_l = 0.0f32; - let mut peak_r = 0.0f32; - let mut true_peak_l = 0.0f32; - let mut true_peak_r = 0.0f32; let mut vu_l_amp = 0.0f32; let mut vu_r_amp = 0.0f32; let frames = interleaved.len() / 2; - let mut wave_l = Vec::with_capacity(frames); - let mut wave_r = Vec::with_capacity(frames); ensure_wave_env_state(&mut ppm_state.wave_env, sample_rate); ensure_lufs_state(&mut ppm_state.lufs, sample_rate); @@ -1218,8 +1239,6 @@ fn build_meter_frame( l = mono; r = mono; } - wave_l.push(l); - wave_r.push(r); wave_env_accumulate(&mut ppm_state.wave_env, l, r, 2); process_lufs_sample(&mut ppm_state.lufs, l, r, rta_config); @@ -1229,10 +1248,11 @@ fn build_meter_frame( ppm_state.din_ppm.process(l, r); ppm_state.ebu_ppm.process(l, r); (vu_l_amp, vu_r_amp) = ppm_state.vu_meter.process(l, r); - peak_l = peak_l.max(abs_l); - peak_r = peak_r.max(abs_r); - true_peak_l = true_peak_l.max(ppm_state.true_peak_l.process(l)); - true_peak_r = true_peak_r.max(ppm_state.true_peak_r.process(r)); + let true_peak_l = ppm_state.true_peak_l.process(l); + let true_peak_r = ppm_state.true_peak_r.process(r); + ppm_state + .transport_peaks + .observe(abs_l, abs_r, true_peak_l, true_peak_r); ppm_state.correlation.process(l, r); ppm_state.xy_pending_l.push(l); @@ -1253,35 +1273,51 @@ fn build_meter_frame( match state { RtaEngineState::Iir(bank) => { finalize_rta_bank(bank, sample_rate, frames, rta_config); - ppm_state.last_rta = Some(build_iir_rta_frame(bank, sample_rate, rta_config)); } RtaEngineState::Fft(fft) => { - if finalize_fft_state(fft, sample_rate, rta_config) { - ppm_state.last_rta = Some(build_fft_rta_frame(fft, sample_rate, rta_config)); - } + let _ = finalize_fft_state(fft, sample_rate, rta_config); } } } - let mut spectro_frame = None; - if let Some(state) = ppm_state.spectro_state.as_mut() { + let spectro_frame = if let Some(state) = ppm_state.spectro_state.as_mut() { if spectro_requested && finalize_spectro_state(state, sample_rate) { - let frame = build_spectro_frame(state, sample_rate); - ppm_state.last_spectro = Some(frame.clone()); - spectro_frame = Some(frame); - } else if !spectro_requested && state.ring_fill >= state.fft_size { - // Keep the current audio window warm, but do not accumulate a - // transform backlog while no client displays the spectrogram. - state.samples_since = state.fft_step_samples; + Some(build_spectro_frame(state, sample_rate)) + } else { + if !spectro_requested && state.ring_fill >= state.fft_size { + // Keep the current audio window warm, but do not accumulate a + // transform backlog while no client displays the spectrogram. + state.samples_since = state.fft_step_samples; + } + None } + } else { + 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 => { + Some(build_fft_rta_frame(fft, sample_rate, rta_config)) + } + RtaEngineState::Fft(_) => None, + }); + + let (transport_peak_l, transport_peak_r) = ppm_state.transport_peaks.take(); + let tp_l = dbfs(transport_peak_l); + let tp_r = dbfs(transport_peak_r); + let rms_l = dbfs(rms_power_l.max(0.0).sqrt()); let rms_r = dbfs(rms_power_r.max(0.0).sqrt()); let vu_l = dbfs(vu_l_amp); let vu_r = dbfs(vu_r_amp); update_box_meter(&mut ppm_state.lufs); - let tp_l = dbfs(true_peak_l.max(peak_l)); - let tp_r = dbfs(true_peak_r.max(peak_r)); let (ppm_din_amp_l, ppm_din_amp_r) = ppm_state.din_ppm.levels(); let (ppm_ebu_amp_l, ppm_ebu_amp_r) = ppm_state.ebu_ppm.levels(); let ppm_din_l = dbfs(ppm_din_amp_l); @@ -1289,17 +1325,9 @@ fn build_meter_frame( let ppm_ebu_l = dbfs(ppm_ebu_amp_l); let ppm_ebu_r = dbfs(ppm_ebu_amp_r); let wave_env = wave_env_flush(&mut ppm_state.wave_env); - let (xy_l, xy_r) = - if ppm_state - .xy_clock - .advance(frames, sample_rate, XY_TARGET_UPDATES_PER_SECOND as u32) - { - take_goniometer_samples(ppm_state, rta_config.xy_points as usize) - } else { - (Vec::new(), Vec::new()) - }; + let (xy_l, xy_r) = take_goniometer_samples(ppm_state, rta_config.xy_points as usize); - MeterFrame { + let frame = MeterFrame { seq: seq.fetch_add(1, Ordering::Relaxed) + 1, timestamp_ms: SystemTime::now() .duration_since(UNIX_EPOCH) @@ -1336,13 +1364,14 @@ fn build_meter_frame( wave_channels: 0, xy_l, xy_r, - rta: ppm_state.last_rta.clone(), - spectro: spectro_frame, + rta, + spectro: None, wave_env, global_config_rev, input, source: "alsa-capture", - } + }; + (Some(frame), spectro_frame) } #[cfg(target_os = "linux")] @@ -1436,7 +1465,6 @@ fn ensure_spectro_state(state: &mut PpmState, sample_rate: u32, config: &Phoenix } state.spectro_signature = signature; state.spectro_state = Some(create_spectro_state(sample_rate, config)); - state.last_spectro = None; } #[cfg(target_os = "linux")] @@ -2226,6 +2254,22 @@ fn band_weighting_gain(f_lo: f32, center: f32, f_hi: f32, mode: &str) -> f32 { mod tests { use super::*; + #[test] + fn transport_peaks_survive_until_the_next_snapshot() { + let mut peaks = TransportPeaks::default(); + peaks.observe(0.1, 0.2, 0.15, 0.25); + peaks.observe(0.8, 0.3, 1.05, 0.45); + peaks.observe(0.2, 0.4, 0.35, 0.5); + + let (left, right) = peaks.take(); + assert_eq!(left, 1.05); + assert_eq!(right, 0.5); + + let (left_after_reset, right_after_reset) = peaks.take(); + assert_eq!(left_after_reset, 0.0); + assert_eq!(right_after_reset, 0.0); + } + fn integrated_fft_power(fft_size: usize) -> f32 { let cycles = 64.0f32; let signal: Vec = (0..fft_size) diff --git a/src/routes.rs b/src/routes.rs index df50b5f..2d75678 100644 --- a/src/routes.rs +++ b/src/routes.rs @@ -890,15 +890,14 @@ 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 mut ticker = tokio::time::interval(std::time::Duration::from_millis(16)); - 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; }; - drain_latest_meter_frame(&mut rx, &mut latest); + let (tp_l, tp_r) = drain_latest_meter_frame(&mut rx, &mut latest); let mut frame = (*latest).clone(); + frame.tp_l = tp_l; + frame.tp_r = tp_r; strip_visual_payloads(&mut frame); let payload = match serde_json::to_string(&frame) { @@ -931,10 +930,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 mut ticker = tokio::time::interval(std::time::Duration::from_millis(16)); - 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; }; @@ -965,34 +961,24 @@ async fn spectro_ws_inner(mut socket: axum::extract::ws::WebSocket, state: AppSt } async fn spectro_ws_session(socket: &mut axum::extract::ws::WebSocket, state: &AppState) { - let mut rx = state.subscribe_metrics(); + let mut rx = state.subscribe_spectro(); loop { - let mut latest = loop { - match rx.recv().await { - Ok(frame) => { - if frame.spectro.is_some() { - break frame; - } - } - Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue, - Err(tokio::sync::broadcast::error::RecvError::Closed) => return, - } + let mut latest = match rx.recv().await { + Ok(frame) => frame, + Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue, + Err(tokio::sync::broadcast::error::RecvError::Closed) => return, }; loop { match rx.try_recv() { - Ok(frame) => { - if frame.spectro.is_some() { - latest = frame; - } - } + Ok(frame) => latest = frame, Err(tokio::sync::broadcast::error::TryRecvError::Lagged(_)) => continue, Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break, Err(tokio::sync::broadcast::error::TryRecvError::Closed) => return, } } - let payload = encode_spectro_frame(latest.spectro.as_ref().expect("checked above")); + let payload = encode_spectro_frame(&latest); if socket .send(axum::extract::ws::Message::Binary(payload)) .await @@ -1099,10 +1085,16 @@ async fn recv_latest_meter_frame( fn drain_latest_meter_frame( rx: &mut tokio::sync::broadcast::Receiver>, latest: &mut Arc, -) { +) -> (f32, f32) { + let mut tp_l = latest.tp_l; + let mut tp_r = latest.tp_r; loop { match rx.try_recv() { - Ok(newer) => *latest = newer, + Ok(newer) => { + tp_l = tp_l.max(newer.tp_l); + tp_r = tp_r.max(newer.tp_r); + *latest = newer; + } Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break, Err(tokio::sync::broadcast::error::TryRecvError::Lagged(skipped)) => { warn!("metrics drain skipped {} stale frames", skipped); @@ -1110,6 +1102,7 @@ fn drain_latest_meter_frame( Err(tokio::sync::broadcast::error::TryRecvError::Closed) => break, } } + (tp_l, tp_r) } fn drain_visual_meter_frames( @@ -1359,6 +1352,33 @@ mod tests { } } + #[test] + fn metrics_drain_keeps_peak_maxima_while_selecting_latest_state() { + let (tx, mut rx) = tokio::sync::broadcast::channel(8); + let mut first = meter_frame(); + first.seq = 1; + first.tp_l = -12.0; + first.tp_r = -9.0; + let mut peak = meter_frame(); + peak.seq = 2; + peak.tp_l = -1.5; + peak.tp_r = -3.0; + let mut latest_frame = meter_frame(); + latest_frame.seq = 3; + latest_frame.tp_l = -18.0; + latest_frame.tp_r = -20.0; + + tx.send(Arc::new(first)).unwrap(); + tx.send(Arc::new(peak)).unwrap(); + tx.send(Arc::new(latest_frame)).unwrap(); + + let mut latest = rx.try_recv().unwrap(); + let (tp_l, tp_r) = drain_latest_meter_frame(&mut rx, &mut latest); + assert_eq!(latest.seq, 3); + assert_eq!(tp_l, -1.5); + assert_eq!(tp_r, -3.0); + } + #[test] fn wave_envelopes_merge_without_losing_columns() { let mut target = Some(WaveEnvFrame { diff --git a/src/state.rs b/src/state.rs index a00bb36..d374473 100644 --- a/src/state.rs +++ b/src/state.rs @@ -17,7 +17,7 @@ use crate::{ config::PhoenixConfig, model::{ InputSource, MeterFrame, PhoenixGlobalConfig, PhoenixGlobalConfigEnvelope, - PhoenixRtaConfig, ServiceStatus, + PhoenixRtaConfig, ServiceStatus, SpectroFrame, }, }; @@ -28,6 +28,7 @@ pub struct AppState { seq: Arc, actual_sample_rate: Arc, metrics_tx: broadcast::Sender>, + spectro_tx: broadcast::Sender>, spectro_subscribers: Arc, restart_token: Arc, rta_config: Arc>, @@ -96,6 +97,7 @@ impl AppState { // WebSocket consumers always drain to the newest state. A small // channel bounds memory and prevents seconds of stale measurements. let (metrics_tx, _) = broadcast::channel(32); + let (spectro_tx, _) = broadcast::channel(8); let initial_global_config = load_global_config(&config); let initial_rta_config = apply_global_to_rta(config.default_rta_config(), &initial_global_config); @@ -106,6 +108,7 @@ impl AppState { seq: Arc::new(AtomicU64::new(0)), actual_sample_rate: Arc::new(AtomicU64::new(configured_sample_rate as u64)), metrics_tx, + spectro_tx, spectro_subscribers: Arc::new(AtomicU64::new(0)), restart_token: Arc::new(AtomicU64::new(0)), rta_config: Arc::new(RwLock::new(initial_rta_config)), @@ -119,6 +122,10 @@ impl AppState { self.metrics_tx.subscribe() } + pub fn subscribe_spectro(&self) -> broadcast::Receiver> { + self.spectro_tx.subscribe() + } + pub fn spectro_subscriber_connected(&self) { self.spectro_subscribers.fetch_add(1, Ordering::SeqCst); } @@ -282,6 +289,7 @@ impl AppState { seq: self.seq.clone(), actual_sample_rate: self.actual_sample_rate.clone(), metrics_tx: self.metrics_tx.clone(), + spectro_tx: self.spectro_tx.clone(), spectro_subscribers: self.spectro_subscribers.clone(), restart_token: self.restart_token.clone(), });