Optimize realtime metrics snapshots
This commit is contained in:
@@ -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.
|
||||
|
||||
+94
-50
@@ -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<AtomicU64>,
|
||||
pub metrics_tx: tokio::sync::broadcast::Sender<Arc<MeterFrame>>,
|
||||
#[cfg_attr(not(target_os = "linux"), allow(dead_code))]
|
||||
pub spectro_tx: tokio::sync::broadcast::Sender<Arc<crate::model::SpectroFrame>>,
|
||||
#[cfg_attr(not(target_os = "linux"), allow(dead_code))]
|
||||
pub spectro_subscribers: Arc<AtomicU64>,
|
||||
#[cfg_attr(not(target_os = "linux"), allow(dead_code))]
|
||||
pub restart_token: Arc<AtomicU64>,
|
||||
@@ -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,8 +553,13 @@ fn capture_until_restart(deps: AudioWorkerDeps, generation: u64) -> anyhow::Resu
|
||||
&rta_config,
|
||||
deps.spectro_subscribers.load(Ordering::Relaxed) > 0,
|
||||
);
|
||||
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 => {
|
||||
warn!("Phoenix ALSA xrun/suspend, preparing capture device again");
|
||||
@@ -550,8 +579,6 @@ struct PpmState {
|
||||
ebu_ppm: PpmDetector,
|
||||
vu_meter: VuMeter,
|
||||
rms_window: MovingAverageWindow,
|
||||
last_rta: Option<RtaFrame>,
|
||||
last_spectro: Option<SpectroFrame>,
|
||||
rta_signature: String,
|
||||
rta_state: Option<RtaEngineState>,
|
||||
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<f32>,
|
||||
xy_pending_r: Vec<f32>,
|
||||
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<f
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
fn build_meter_frame(
|
||||
fn process_audio_block(
|
||||
interleaved: &[i16],
|
||||
sample_rate: u32,
|
||||
input: InputSource,
|
||||
@@ -1170,18 +1197,12 @@ fn build_meter_frame(
|
||||
ppm_state: &mut PpmState,
|
||||
rta_config: &PhoenixRtaConfig,
|
||||
spectro_requested: bool,
|
||||
) -> MeterFrame {
|
||||
) -> (Option<MeterFrame>, Option<SpectroFrame>) {
|
||||
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 {
|
||||
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<f32> = (0..fft_size)
|
||||
|
||||
+44
-24
@@ -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;
|
||||
}
|
||||
}
|
||||
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<Arc<MeterFrame>>,
|
||||
latest: &mut Arc<MeterFrame>,
|
||||
) {
|
||||
) -> (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 {
|
||||
|
||||
+9
-1
@@ -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<AtomicU64>,
|
||||
actual_sample_rate: Arc<AtomicU64>,
|
||||
metrics_tx: broadcast::Sender<Arc<MeterFrame>>,
|
||||
spectro_tx: broadcast::Sender<Arc<SpectroFrame>>,
|
||||
spectro_subscribers: Arc<AtomicU64>,
|
||||
restart_token: Arc<AtomicU64>,
|
||||
rta_config: Arc<RwLock<PhoenixRtaConfig>>,
|
||||
@@ -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<Arc<SpectroFrame>> {
|
||||
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(),
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user