Optimize realtime display processing

This commit is contained in:
Mikei386 committed 2026-07-21 22:04:50 +02:00
1 parent a123539023
commit c067c73024
6 files changed
+165 -26

No files matched your search

+9 -1
View File
@@ -102,6 +102,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_subscribers: Arc<AtomicU64>,
#[cfg_attr(not(target_os = "linux"), allow(dead_code))]
pub restart_token: Arc<AtomicU64>,
}
@@ -437,6 +439,7 @@ fn capture_until_restart(deps: AudioWorkerDeps, generation: u64) -> anyhow::Resu
&deps.seq,
&mut ppm,
&rta_config,
deps.spectro_subscribers.load(Ordering::Relaxed) > 0,
);
let _ = deps.metrics_tx.send(Arc::new(frame));
}
@@ -1078,6 +1081,7 @@ fn build_meter_frame(
seq: &Arc<AtomicU64>,
ppm_state: &mut PpmState,
rta_config: &PhoenixRtaConfig,
spectro_requested: bool,
) -> MeterFrame {
let mut rms_power_l = 0.0f32;
let mut rms_power_r = 0.0f32;
@@ -1172,10 +1176,14 @@ fn build_meter_frame(
}
let mut spectro_frame = None;
if let Some(state) = ppm_state.spectro_state.as_mut() {
if finalize_spectro_state(state, sample_rate) {
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;
}
}
+6
View File
@@ -959,6 +959,12 @@ async fn visuals_ws_inner(mut socket: axum::extract::ws::WebSocket, state: AppSt
}
async fn spectro_ws_inner(mut socket: axum::extract::ws::WebSocket, state: AppState) {
state.spectro_subscriber_connected();
spectro_ws_session(&mut socket, &state).await;
state.spectro_subscriber_disconnected();
}
async fn spectro_ws_session(socket: &mut axum::extract::ws::WebSocket, state: &AppState) {
let mut rx = state.subscribe_metrics();
loop {
let mut latest = loop {
+15
View File
@@ -28,6 +28,7 @@ pub struct AppState {
seq: Arc<AtomicU64>,
actual_sample_rate: Arc<AtomicU64>,
metrics_tx: broadcast::Sender<Arc<MeterFrame>>,
spectro_subscribers: Arc<AtomicU64>,
restart_token: Arc<AtomicU64>,
rta_config: Arc<RwLock<PhoenixRtaConfig>>,
global_config: Arc<RwLock<PhoenixGlobalConfig>>,
@@ -105,6 +106,7 @@ impl AppState {
seq: Arc::new(AtomicU64::new(0)),
actual_sample_rate: Arc::new(AtomicU64::new(configured_sample_rate as u64)),
metrics_tx,
spectro_subscribers: Arc::new(AtomicU64::new(0)),
restart_token: Arc::new(AtomicU64::new(0)),
rta_config: Arc::new(RwLock::new(initial_rta_config)),
global_config: Arc::new(RwLock::new(initial_global_config)),
@@ -117,6 +119,18 @@ impl AppState {
self.metrics_tx.subscribe()
}
pub fn spectro_subscriber_connected(&self) {
self.spectro_subscribers.fetch_add(1, Ordering::SeqCst);
}
pub fn spectro_subscriber_disconnected(&self) {
let _ =
self.spectro_subscribers
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |count| {
Some(count.saturating_sub(1))
});
}
pub async fn input(&self) -> InputSource {
*self.current_input.read().await
}
@@ -268,6 +282,7 @@ impl AppState {
seq: self.seq.clone(),
actual_sample_rate: self.actual_sample_rate.clone(),
metrics_tx: self.metrics_tx.clone(),
spectro_subscribers: self.spectro_subscribers.clone(),
restart_token: self.restart_token.clone(),
});
}