Revert "Decouple realtime snapshots from transport cadence"
This reverts commit feec2bac56.
This commit is contained in:
1 parent
dba83a2d98
commit
96d323c32d
3 files changed
+52
-82
No files matched your search
+17
-4
@@ -14,7 +14,7 @@ use tracing::warn;
|
|||||||
#[cfg(target_os = "linux")]
|
#[cfg(target_os = "linux")]
|
||||||
use crate::correlation::CorrelationMeter;
|
use crate::correlation::CorrelationMeter;
|
||||||
#[cfg(target_os = "linux")]
|
#[cfg(target_os = "linux")]
|
||||||
use crate::goniometer::selected_sample_indices;
|
use crate::goniometer::{selected_sample_indices, GoniometerClock};
|
||||||
#[cfg(target_os = "linux")]
|
#[cfg(target_os = "linux")]
|
||||||
use crate::model::{RtaFrame, SpectroFrame, WaveEnvFrame};
|
use crate::model::{RtaFrame, SpectroFrame, WaveEnvFrame};
|
||||||
#[cfg(target_os = "linux")]
|
#[cfg(target_os = "linux")]
|
||||||
@@ -40,8 +40,12 @@ use crate::{
|
|||||||
model::{InputSource, MeterFrame, PhoenixRtaConfig},
|
model::{InputSource, MeterFrame, PhoenixRtaConfig},
|
||||||
};
|
};
|
||||||
|
|
||||||
#[cfg(not(target_os = "linux"))]
|
// Keep transport sampling safely above the 60 Hz display cadence. ALSA only
|
||||||
const PLACEHOLDER_UPDATES_PER_SECOND: u32 = 120;
|
// 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)]
|
#[derive(Clone, Copy, Debug, Default)]
|
||||||
struct TransportPeaks {
|
struct TransportPeaks {
|
||||||
@@ -415,7 +419,7 @@ pub fn spawn_audio_capture_worker(deps: AudioWorkerDeps) {
|
|||||||
pub fn spawn_audio_capture_worker(deps: AudioWorkerDeps) {
|
pub fn spawn_audio_capture_worker(deps: AudioWorkerDeps) {
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
warn!("Phoenix ALSA capture is only available on Linux; emitting placeholder frames on this host");
|
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));
|
let mut ticker = tokio::time::interval(Duration::from_nanos(tick_ns));
|
||||||
deps.actual_sample_rate
|
deps.actual_sample_rate
|
||||||
.store(deps.config.sample_rate as u64, Ordering::SeqCst);
|
.store(deps.config.sample_rate as u64, Ordering::SeqCst);
|
||||||
@@ -594,6 +598,7 @@ struct PpmState {
|
|||||||
lr_delay_y1_r: f32,
|
lr_delay_y1_r: f32,
|
||||||
true_peak_l: TruePeakDetector,
|
true_peak_l: TruePeakDetector,
|
||||||
true_peak_r: TruePeakDetector,
|
true_peak_r: TruePeakDetector,
|
||||||
|
transport_clock: GoniometerClock,
|
||||||
transport_peaks: TransportPeaks,
|
transport_peaks: TransportPeaks,
|
||||||
correlation: CorrelationMeter,
|
correlation: CorrelationMeter,
|
||||||
phase_wheel: PhaseWheelAnalyzer,
|
phase_wheel: PhaseWheelAnalyzer,
|
||||||
@@ -626,6 +631,7 @@ impl Default for PpmState {
|
|||||||
lr_delay_y1_r: 0.0,
|
lr_delay_y1_r: 0.0,
|
||||||
true_peak_l: TruePeakDetector::default(),
|
true_peak_l: TruePeakDetector::default(),
|
||||||
true_peak_r: TruePeakDetector::default(),
|
true_peak_r: TruePeakDetector::default(),
|
||||||
|
transport_clock: GoniometerClock::default(),
|
||||||
transport_peaks: TransportPeaks::default(),
|
transport_peaks: TransportPeaks::default(),
|
||||||
correlation: CorrelationMeter::new(48_000, 1.0, -75.0, 0),
|
correlation: CorrelationMeter::new(48_000, 1.0, -75.0, 0),
|
||||||
phase_wheel: PhaseWheelAnalyzer::new(48_000),
|
phase_wheel: PhaseWheelAnalyzer::new(48_000),
|
||||||
@@ -1251,6 +1257,13 @@ fn process_audio_block(
|
|||||||
None
|
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 {
|
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::Iir(bank) => Some(build_iir_rta_frame(bank, sample_rate, rta_config)),
|
||||||
RtaEngineState::Fft(fft) if fft.ring_fill >= fft.fft_size => {
|
RtaEngineState::Fft(fft) if fft.ring_fill >= fft.fft_size => {
|
||||||
|
|||||||
+25
-2
@@ -1,12 +1,10 @@
|
|||||||
//! Timing and sample-selection helpers for the realtime goniometer stream.
|
//! Timing and sample-selection helpers for the realtime goniometer stream.
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
#[derive(Clone, Debug, Default)]
|
#[derive(Clone, Debug, Default)]
|
||||||
pub struct GoniometerClock {
|
pub struct GoniometerClock {
|
||||||
phase: u64,
|
phase: u64,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
impl GoniometerClock {
|
impl GoniometerClock {
|
||||||
/// Returns true at the first capture boundary after the next display tick.
|
/// Returns true at the first capture boundary after the next display tick.
|
||||||
/// The fractional phase is retained, so the average rate is independent
|
/// 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]
|
#[test]
|
||||||
fn selection_never_upsamples_and_keeps_endpoints() {
|
fn selection_never_upsamples_and_keeps_endpoints() {
|
||||||
assert_eq!(selected_sample_indices(3, 1024), vec![0, 1, 2]);
|
assert_eq!(selected_sample_indices(3, 1024), vec![0, 1, 2]);
|
||||||
|
|||||||
+10
-76
@@ -14,7 +14,7 @@ use std::{
|
|||||||
process::Stdio,
|
process::Stdio,
|
||||||
sync::atomic::{AtomicU64, Ordering},
|
sync::atomic::{AtomicU64, Ordering},
|
||||||
sync::Arc,
|
sync::Arc,
|
||||||
time::{Duration, SystemTime, UNIX_EPOCH},
|
time::{SystemTime, UNIX_EPOCH},
|
||||||
};
|
};
|
||||||
use tokio::{fs, io::AsyncWriteExt, process::Command};
|
use tokio::{fs, io::AsyncWriteExt, process::Command};
|
||||||
use tracing::warn;
|
use tracing::warn;
|
||||||
@@ -27,8 +27,6 @@ use crate::{
|
|||||||
const ONLINE_UPDATE_ZIP_URL: &str = "https://webshare.casaderoll.de/share/Phoenix.zip";
|
const ONLINE_UPDATE_ZIP_URL: &str = "https://webshare.casaderoll.de/share/Phoenix.zip";
|
||||||
const DEVICE_BACKUP_KIND: &str = "phoenix-device-backup";
|
const DEVICE_BACKUP_KIND: &str = "phoenix-device-backup";
|
||||||
const DEVICE_BACKUP_SCHEMA_VERSION: u32 = 1;
|
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);
|
static JSON_WRITE_TOKEN: AtomicU64 = AtomicU64::new(1);
|
||||||
|
|
||||||
fn stable_script_command(program: &str) -> Command {
|
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) {
|
async fn metrics_ws_inner(mut socket: axum::extract::ws::WebSocket, state: AppState) {
|
||||||
let mut rx = state.subscribe_metrics();
|
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 {
|
loop {
|
||||||
ticker.tick().await;
|
|
||||||
let Some(mut latest) = recv_latest_meter_frame(&mut rx).await else {
|
let Some(mut latest) = recv_latest_meter_frame(&mut rx).await else {
|
||||||
break;
|
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) {
|
async fn visuals_ws_inner(mut socket: axum::extract::ws::WebSocket, state: AppState) {
|
||||||
let mut rx = state.subscribe_metrics();
|
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 {
|
loop {
|
||||||
ticker.tick().await;
|
|
||||||
let Some(mut latest) = recv_latest_meter_frame(&mut rx).await else {
|
let Some(mut latest) = recv_latest_meter_frame(&mut rx).await else {
|
||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
@@ -1426,14 +1416,19 @@ fn drain_visual_meter_frames(
|
|||||||
latest: &mut Arc<MeterFrame>,
|
latest: &mut Arc<MeterFrame>,
|
||||||
) -> (Option<WaveEnvFrame>, Option<(Vec<f32>, Vec<f32>)>) {
|
) -> (Option<WaveEnvFrame>, Option<(Vec<f32>, Vec<f32>)>) {
|
||||||
let mut combined_wave_env = latest.wave_env.clone();
|
let mut combined_wave_env = latest.wave_env.clone();
|
||||||
let mut combined_xy = None;
|
let mut latest_xy = if latest.xy_l.is_empty() || latest.xy_r.is_empty() {
|
||||||
append_xy_samples(&mut combined_xy, &latest.xy_l, &latest.xy_r);
|
None
|
||||||
|
} else {
|
||||||
|
Some((latest.xy_l.clone(), latest.xy_r.clone()))
|
||||||
|
};
|
||||||
|
|
||||||
loop {
|
loop {
|
||||||
match rx.try_recv() {
|
match rx.try_recv() {
|
||||||
Ok(newer) => {
|
Ok(newer) => {
|
||||||
merge_wave_env(&mut combined_wave_env, newer.wave_env.clone());
|
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;
|
*latest = newer;
|
||||||
}
|
}
|
||||||
Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break,
|
Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break,
|
||||||
@@ -1444,31 +1439,7 @@ fn drain_visual_meter_frames(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
(combined_wave_env, combined_xy)
|
(combined_wave_env, latest_xy)
|
||||||
}
|
|
||||||
|
|
||||||
fn append_xy_samples(
|
|
||||||
target: &mut Option<(Vec<f32>, Vec<f32>)>,
|
|
||||||
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);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn merge_wave_env(target: &mut Option<WaveEnvFrame>, incoming: Option<WaveEnvFrame>) {
|
fn merge_wave_env(target: &mut Option<WaveEnvFrame>, incoming: Option<WaveEnvFrame>) {
|
||||||
@@ -1718,43 +1689,6 @@ mod tests {
|
|||||||
assert_eq!(tp_r, -3.0);
|
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<f32> = (0..1500).map(|value| value as f32).collect();
|
|
||||||
let second: Vec<f32> = (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]
|
#[test]
|
||||||
fn wave_envelopes_merge_without_losing_columns() {
|
fn wave_envelopes_merge_without_losing_columns() {
|
||||||
let mut target = Some(WaveEnvFrame {
|
let mut target = Some(WaveEnvFrame {
|
||||||
|
|||||||
Reference in new issue
Block a user