From fd69cb49885f7ec0544ac94f48bb02bfe7fbaaf9 Mon Sep 17 00:00:00 2001 From: Psychotoxical <171614930+Psychotoxical@users.noreply.github.com> Date: Fri, 8 May 2026 12:20:40 +0200 Subject: [PATCH] refactor(lib): extract analysis runtime queues into own module MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Move ~440 LOC of analysis-queue plumbing out of lib.rs into a new src-tauri/src/analysis_runtime.rs: - AnalysisBackfillQueueState/Shared + worker loop + lazy-init - AnalysisCpuSeedQueueState/Shared + worker loop + lazy-init - submit_analysis_cpu_seed - emit_analysis_queue_snapshot_line + analysis_queue_snapshot_loop - WaveformUpdatedPayload (only used internally) External callers keep their existing paths via `pub(crate) use analysis_runtime::*` re-export at the lib.rs root. Direct accesses to the static globals from app_api/analysis.rs are replaced with a new prune_analysis_queues(keep) function so the statics stay private. lib.rs goes from 999 → ~540 LOC. sync_cancel_flags stays put (sync domain, separate concern). Tray-related types stay too — they're a separate split candidate. Behaviour-preserving: same workers, same queues, same events, same return shapes. cargo check clean. --- src-tauri/src/analysis_runtime.rs | 496 ++++++++++++++++++ src-tauri/src/lib.rs | 458 +--------------- .../src/lib_commands/app_api/analysis.rs | 21 +- 3 files changed, 502 insertions(+), 473 deletions(-) create mode 100644 src-tauri/src/analysis_runtime.rs diff --git a/src-tauri/src/analysis_runtime.rs b/src-tauri/src/analysis_runtime.rs new file mode 100644 index 00000000..afa262c3 --- /dev/null +++ b/src-tauri/src/analysis_runtime.rs @@ -0,0 +1,496 @@ +use std::collections::{HashSet, VecDeque}; +use std::sync::{Arc, Mutex, OnceLock}; + +use tauri::{Emitter, Manager}; + +use super::analysis_cache; +use super::lib_commands::*; +use super::subsonic_wire_user_agent; + +#[derive(Clone, serde::Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct WaveformUpdatedPayload { + pub(crate) track_id: String, + pub(crate) is_partial: bool, +} + +// ─── HTTP backfill queue: download tracks + seed analysis cache ────────────── + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum AnalysisBackfillEnqueueKind { + /// New job at the tail of the queue. + NewBack, + /// New job for the currently playing track (head). + NewFront, + /// Same track was already waiting; moved to head with the latest URL. + ReorderedFront, + /// Low-priority duplicate while the track is already queued or running. + DuplicateSkipped, + /// High-priority request but that track is already being downloaded+seeded. + RunningSkipped, +} + +#[derive(Default)] +pub(crate) struct AnalysisBackfillQueueState { + pub(crate) deque: VecDeque<(String, String)>, + /// Set while this `track_id` is inside `analysis_backfill_download_and_seed` (not in deque). + pub(crate) in_progress: Option, +} + +impl AnalysisBackfillQueueState { + fn is_reserved(&self, tid: &str) -> bool { + self.in_progress.as_deref() == Some(tid) + || self.deque.iter().any(|(t, _)| t.as_str() == tid) + } + + fn try_pop_next(&mut self) -> Option<(String, String)> { + let (tid, url) = self.deque.pop_front()?; + self.in_progress = Some(tid.clone()); + Some((tid, url)) + } + + fn finish_job(&mut self, tid: &str) { + if self.in_progress.as_deref() == Some(tid) { + self.in_progress = None; + } + } + + pub(crate) fn enqueue( + &mut self, + tid: String, + url: String, + high_priority: bool, + ) -> AnalysisBackfillEnqueueKind { + let tref = tid.as_str(); + if self.is_reserved(tref) { + if !high_priority { + return AnalysisBackfillEnqueueKind::DuplicateSkipped; + } + if self.in_progress.as_deref() == Some(tref) { + return AnalysisBackfillEnqueueKind::RunningSkipped; + } + self.deque.retain(|(t, _)| t != &tid); + self.deque.push_front((tid, url)); + return AnalysisBackfillEnqueueKind::ReorderedFront; + } + if high_priority { + self.deque.push_front((tid, url)); + AnalysisBackfillEnqueueKind::NewFront + } else { + self.deque.push_back((tid, url)); + AnalysisBackfillEnqueueKind::NewBack + } + } + + pub(crate) fn prune_queued_not_in(&mut self, keep_track_ids: &HashSet<&str>) -> usize { + let before = self.deque.len(); + self.deque + .retain(|(track_id, _)| keep_track_ids.contains(track_id.as_str())); + before.saturating_sub(self.deque.len()) + } +} + +pub(crate) struct AnalysisBackfillShared { + pub(crate) state: Mutex, + wake_tx: tokio::sync::mpsc::UnboundedSender<()>, +} + +impl AnalysisBackfillShared { + pub(crate) fn ping_worker(&self) { + let _ = self.wake_tx.send(()); + } +} + +static ANALYSIS_BACKFILL: OnceLock> = OnceLock::new(); + +/// Lazily spawns the single backfill worker (first caller supplies `AppHandle`). +pub(crate) fn analysis_backfill_shared(app: &tauri::AppHandle) -> Arc { + ANALYSIS_BACKFILL + .get_or_init(|| { + let (wake_tx, wake_rx) = tokio::sync::mpsc::unbounded_channel(); + let shared = Arc::new(AnalysisBackfillShared { + state: Mutex::new(AnalysisBackfillQueueState::default()), + wake_tx, + }); + let app = app.clone(); + let sh = shared.clone(); + tauri::async_runtime::spawn(analysis_backfill_worker_loop(app, sh, wake_rx)); + shared + }) + .clone() +} + +async fn analysis_backfill_download_and_seed( + app: &tauri::AppHandle, + track_id: &str, + url: &str, +) -> Result { + let client = reqwest::Client::builder() + .user_agent(subsonic_wire_user_agent()) + .timeout(std::time::Duration::from_secs(120)) + .build() + .map_err(|e| e.to_string())?; + let response = client.get(url).send().await.map_err(|e| e.to_string())?; + if !response.status().is_success() { + return Err(format!("HTTP {}", response.status().as_u16())); + } + let bytes = response.bytes().await.map_err(|e| e.to_string())?; + if bytes.is_empty() { + return Err("empty response".to_string()); + } + enqueue_analysis_seed(app, track_id, &bytes).await +} + +async fn analysis_backfill_worker_loop( + app: tauri::AppHandle, + shared: Arc, + mut wake_rx: tokio::sync::mpsc::UnboundedReceiver<()>, +) { + loop { + if wake_rx.recv().await.is_none() { + break; + } + while let Some((track_id, url)) = { + let mut st = shared + .state + .lock() + .unwrap_or_else(|e| e.into_inner()); + st.try_pop_next() + } { + crate::app_deprintln!("[analysis] backfill worker: start track_id={}", track_id); + let result = analysis_backfill_download_and_seed(&app, &track_id, &url).await; + match &result { + Ok(has_loudness) => crate::app_deprintln!( + "[analysis] backfill ready: {} (has_loudness={})", + track_id, + has_loudness + ), + Err(e) => crate::app_eprintln!("[analysis] backfill failed for {}: {}", track_id, e), + } + let mut st = shared + .state + .lock() + .unwrap_or_else(|e| e.into_inner()); + st.finish_job(&track_id); + } + } +} + +pub(crate) fn analysis_backfill_is_current_track(app: &tauri::AppHandle, track_id: &str) -> bool { + app.try_state::() + .is_some_and(|e| crate::audio::analysis_track_id_is_current_playback(&e, track_id)) +} + +// ─── Full-track waveform + loudness: single CPU worker (mirrors HTTP backfill queue) ─ +// One `spawn_blocking` decode at a time; current playback is high-priority (front + reorder). +// Same `track_id` queued again merges waiters onto one job; while decode runs, same-id +// submitters attach to `running` followers so they all get the same outcome. + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum AnalysisCpuSeedEnqueueKind { + NewBack, + NewFront, + ReorderedFront, + RunningFollower, + MergedQueued, +} + +struct AnalysisCpuSeedJob { + track_id: String, + bytes: Vec, + waiters: Vec>>, +} + +struct AnalysisCpuSeedQueueState { + deque: VecDeque, + /// Decode in progress — same-id callers wait here for the same outcome. + running: Option<( + String, + Arc>>>>, + )>, +} + +impl AnalysisCpuSeedQueueState { + fn enqueue( + &mut self, + track_id: String, + bytes: Vec, + high_priority: bool, + ) -> ( + AnalysisCpuSeedEnqueueKind, + tokio::sync::oneshot::Receiver>, + ) { + let (done_tx, done_rx) = tokio::sync::oneshot::channel(); + let tid = track_id.as_str(); + + if let Some((rtid, followers)) = &self.running { + if rtid == tid { + followers + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(done_tx); + return (AnalysisCpuSeedEnqueueKind::RunningFollower, done_rx); + } + } + + if let Some(pos) = self.deque.iter().position(|j| j.track_id == track_id) { + let mut job = self.deque.remove(pos).unwrap(); + job.bytes = bytes; + job.waiters.push(done_tx); + let kind = if high_priority { + self.deque.push_front(job); + AnalysisCpuSeedEnqueueKind::ReorderedFront + } else { + self.deque.push_back(job); + AnalysisCpuSeedEnqueueKind::MergedQueued + }; + return (kind, done_rx); + } + + let job = AnalysisCpuSeedJob { + track_id: track_id.clone(), + bytes, + waiters: vec![done_tx], + }; + let kind = if high_priority { + self.deque.push_front(job); + AnalysisCpuSeedEnqueueKind::NewFront + } else { + self.deque.push_back(job); + AnalysisCpuSeedEnqueueKind::NewBack + }; + (kind, done_rx) + } + + fn prune_queued_not_in(&mut self, keep_track_ids: &HashSet<&str>) -> (usize, usize) { + let mut kept = VecDeque::with_capacity(self.deque.len()); + let mut removed_jobs = 0usize; + let mut removed_waiters = 0usize; + while let Some(job) = self.deque.pop_front() { + if keep_track_ids.contains(job.track_id.as_str()) { + kept.push_back(job); + continue; + } + removed_jobs += 1; + removed_waiters += job.waiters.len(); + for tx in job.waiters { + let _ = tx.send(Err( + "cpu-seed pruned: track no longer in playback queue".to_string(), + )); + } + } + self.deque = kept; + (removed_jobs, removed_waiters) + } +} + +struct AnalysisCpuSeedShared { + state: Mutex, + wake_tx: tokio::sync::mpsc::UnboundedSender<()>, +} + +impl Default for AnalysisCpuSeedQueueState { + fn default() -> Self { + Self { + deque: VecDeque::new(), + running: None, + } + } +} + +impl AnalysisCpuSeedShared { + fn ping_worker(&self) { + let _ = self.wake_tx.send(()); + } +} + +static ANALYSIS_CPU_SEED: OnceLock> = OnceLock::new(); + +fn analysis_cpu_seed_shared(app: &tauri::AppHandle) -> Arc { + ANALYSIS_CPU_SEED + .get_or_init(|| { + let (wake_tx, wake_rx) = tokio::sync::mpsc::unbounded_channel(); + let shared = Arc::new(AnalysisCpuSeedShared { + state: Mutex::new(AnalysisCpuSeedQueueState::default()), + wake_tx, + }); + let app = app.clone(); + let sh = shared.clone(); + tauri::async_runtime::spawn(analysis_cpu_seed_worker_loop(app, sh, wake_rx)); + shared + }) + .clone() +} + +/// HTTP backfill + CPU seed queue sizes (debug log only — `app_deprintln!`). +fn emit_analysis_queue_snapshot_line() { + let http = if let Some(arc) = ANALYSIS_BACKFILL.get() { + let st = arc.state.lock().unwrap_or_else(|e| e.into_inner()); + format!( + "http_backfill={{queued:{} download_active:{:?}}}", + st.deque.len(), + st.in_progress.as_deref() + ) + } else { + "http_backfill={{not_started}}".to_string() + }; + + let cpu = if let Some(arc) = ANALYSIS_CPU_SEED.get() { + let st = arc.state.lock().unwrap_or_else(|e| e.into_inner()); + let queued_jobs = st.deque.len(); + let pending_in_queued_jobs: usize = st.deque.iter().map(|j| j.waiters.len()).sum(); + let (decoding_tid, decoding_extra_waiters) = match &st.running { + Some((tid, fl)) => ( + Some(tid.as_str()), + fl.lock().map(|g| g.len()).unwrap_or(0), + ), + None => (None, 0usize), + }; + format!( + "cpu_seed={{queued_jobs:{} pending_channels_in_queue:{} decoding_tid:{:?} extra_waiters_same_id:{}}}", + queued_jobs, + pending_in_queued_jobs, + decoding_tid, + decoding_extra_waiters + ) + } else { + "cpu_seed={{not_started}}".to_string() + }; + + crate::app_deprintln!( + "[analysis] queue_snapshot interval_s=60 note=queues_in_memory_cleared_on_app_restart | {http} | {cpu}" + ); +} + +pub(crate) async fn analysis_queue_snapshot_loop() { + emit_analysis_queue_snapshot_line(); + loop { + tokio::time::sleep(std::time::Duration::from_secs(60)).await; + emit_analysis_queue_snapshot_line(); + } +} + +async fn analysis_cpu_seed_worker_loop( + app: tauri::AppHandle, + shared: Arc, + mut wake_rx: tokio::sync::mpsc::UnboundedReceiver<()>, +) { + loop { + if wake_rx.recv().await.is_none() { + break; + } + loop { + let (job, followers) = { + let mut st = shared.state.lock().unwrap_or_else(|e| e.into_inner()); + let Some(j) = st.deque.pop_front() else { + break; + }; + let fl = Arc::new(Mutex::new(Vec::new())); + st.running = Some((j.track_id.clone(), fl.clone())); + (j, fl) + }; + let tid_log = job.track_id.clone(); + let app2 = app.clone(); + let tid = job.track_id.clone(); + let bytes = job.bytes; + let outcome = tokio::task::spawn_blocking(move || { + analysis_cache::seed_from_bytes_execute(&app2, &tid, &bytes) + }) + .await + .unwrap_or_else(|e| Err(format!("cpu-seed spawn_blocking: {e}"))); + + let mut extra = followers + .lock() + .unwrap_or_else(|e| e.into_inner()) + .drain(..) + .collect::>(); + for tx in job.waiters { + let _ = tx.send(outcome.clone()); + } + for tx in extra.drain(..) { + let _ = tx.send(outcome.clone()); + } + + { + let mut st = shared.state.lock().unwrap_or_else(|e| e.into_inner()); + st.running = None; + } + let ok = outcome.as_ref().map(|o| *o == analysis_cache::SeedFromBytesOutcome::Upserted).unwrap_or(false); + crate::app_deprintln!( + "[analysis] cpu-seed worker: done track_id={} upserted={}", + tid_log, + ok + ); + } + } +} + +/// Prune queued items in both analysis queues (HTTP backfill + CPU seed) whose +/// track ids are not in `keep_track_ids`. Items that are *currently running* are +/// untouched; only queued items are removed. Pruned CPU-seed waiters get an Err +/// indicating the prune. +/// +/// Returns `(http_removed, cpu_removed_jobs, cpu_removed_waiters)`. Either +/// queue may not have been initialized yet — those slots return 0. +pub(crate) fn prune_analysis_queues( + keep_track_ids: &HashSet<&str>, +) -> Result<(usize, usize, usize), String> { + let http_removed = if let Some(shared) = ANALYSIS_BACKFILL.get() { + let mut st = shared + .state + .lock() + .map_err(|_| "analysis backfill lock poisoned".to_string())?; + st.prune_queued_not_in(keep_track_ids) + } else { + 0 + }; + + let (cpu_removed_jobs, cpu_removed_waiters) = if let Some(shared) = ANALYSIS_CPU_SEED.get() { + let mut st = shared + .state + .lock() + .map_err(|_| "analysis cpu-seed lock poisoned".to_string())?; + st.prune_queued_not_in(keep_track_ids) + } else { + (0, 0) + }; + + Ok((http_removed, cpu_removed_jobs, cpu_removed_waiters)) +} + +/// Submit full-buffer analysis; serializes with other producers. `high_priority` mirrors +/// HTTP backfill head insertion for the currently playing track. +/// +/// Emits `analysis:waveform-updated` when analysis **wrote** new waveform data (`Upserted`). +/// Cache-hit skips (`SkippedWaveformCacheHit`) omit the event so the frontend does not +/// re-run loudness refresh / waveform IPC for rows that were already current. +pub(crate) async fn submit_analysis_cpu_seed( + app: tauri::AppHandle, + track_id: String, + bytes: Vec, + high_priority: bool, +) -> Result { + let shared = analysis_cpu_seed_shared(&app); + let rx = { + let mut st = shared.state.lock().unwrap_or_else(|e| e.into_inner()); + let (kind, rx) = st.enqueue(track_id.clone(), bytes, high_priority); + crate::app_deprintln!("[analysis] cpu-seed submit: kind={kind:?} high_priority={high_priority}"); + drop(st); + shared.ping_worker(); + rx + }; + let outcome = match rx.await { + Ok(res) => res?, + Err(_) => return Err("cpu-seed: result channel dropped".to_string()), + }; + if matches!(outcome, analysis_cache::SeedFromBytesOutcome::Upserted) { + let _ = app.emit( + "analysis:waveform-updated", + WaveformUpdatedPayload { + track_id: track_id.clone(), + is_partial: false, + }, + ); + } + Ok(outcome) +} diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index 1796563b..5e753a3e 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -3,6 +3,7 @@ mod audio; mod analysis_cache; +mod analysis_runtime; pub mod cli; mod discord; pub(crate) mod logging; @@ -10,7 +11,9 @@ mod lib_commands; #[cfg(target_os = "windows")] mod taskbar_win; -use std::collections::{HashMap, HashSet, VecDeque}; +pub(crate) use analysis_runtime::*; + +use std::collections::HashMap; use std::sync::{Arc, Mutex, OnceLock, RwLock}; use std::sync::atomic::{AtomicBool, Ordering}; @@ -62,452 +65,6 @@ fn sync_cancel_flags() -> &'static Mutex>> { FLAGS.get_or_init(|| Mutex::new(HashMap::new())) } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -enum AnalysisBackfillEnqueueKind { - /// New job at the tail of the queue. - NewBack, - /// New job for the currently playing track (head). - NewFront, - /// Same track was already waiting; moved to head with the latest URL. - ReorderedFront, - /// Low-priority duplicate while the track is already queued or running. - DuplicateSkipped, - /// High-priority request but that track is already being downloaded+seeded. - RunningSkipped, -} - -#[derive(Default)] -struct AnalysisBackfillQueueState { - deque: VecDeque<(String, String)>, - /// Set while this `track_id` is inside `analysis_backfill_download_and_seed` (not in deque). - in_progress: Option, -} - -impl AnalysisBackfillQueueState { - fn is_reserved(&self, tid: &str) -> bool { - self.in_progress.as_deref() == Some(tid) - || self.deque.iter().any(|(t, _)| t.as_str() == tid) - } - - fn try_pop_next(&mut self) -> Option<(String, String)> { - let (tid, url) = self.deque.pop_front()?; - self.in_progress = Some(tid.clone()); - Some((tid, url)) - } - - fn finish_job(&mut self, tid: &str) { - if self.in_progress.as_deref() == Some(tid) { - self.in_progress = None; - } - } - - fn enqueue( - &mut self, - tid: String, - url: String, - high_priority: bool, - ) -> AnalysisBackfillEnqueueKind { - let tref = tid.as_str(); - if self.is_reserved(tref) { - if !high_priority { - return AnalysisBackfillEnqueueKind::DuplicateSkipped; - } - if self.in_progress.as_deref() == Some(tref) { - return AnalysisBackfillEnqueueKind::RunningSkipped; - } - self.deque.retain(|(t, _)| t != &tid); - self.deque.push_front((tid, url)); - return AnalysisBackfillEnqueueKind::ReorderedFront; - } - if high_priority { - self.deque.push_front((tid, url)); - AnalysisBackfillEnqueueKind::NewFront - } else { - self.deque.push_back((tid, url)); - AnalysisBackfillEnqueueKind::NewBack - } - } - - fn prune_queued_not_in(&mut self, keep_track_ids: &HashSet<&str>) -> usize { - let before = self.deque.len(); - self.deque - .retain(|(track_id, _)| keep_track_ids.contains(track_id.as_str())); - before.saturating_sub(self.deque.len()) - } -} - -struct AnalysisBackfillShared { - state: Mutex, - wake_tx: tokio::sync::mpsc::UnboundedSender<()>, -} - -impl AnalysisBackfillShared { - fn ping_worker(&self) { - let _ = self.wake_tx.send(()); - } -} - -static ANALYSIS_BACKFILL: OnceLock> = OnceLock::new(); - -/// Lazily spawns the single backfill worker (first caller supplies `AppHandle`). -fn analysis_backfill_shared(app: &tauri::AppHandle) -> Arc { - ANALYSIS_BACKFILL - .get_or_init(|| { - let (wake_tx, wake_rx) = tokio::sync::mpsc::unbounded_channel(); - let shared = Arc::new(AnalysisBackfillShared { - state: Mutex::new(AnalysisBackfillQueueState::default()), - wake_tx, - }); - let app = app.clone(); - let sh = shared.clone(); - tauri::async_runtime::spawn(analysis_backfill_worker_loop(app, sh, wake_rx)); - shared - }) - .clone() -} - -async fn analysis_backfill_download_and_seed( - app: &tauri::AppHandle, - track_id: &str, - url: &str, -) -> Result { - let client = reqwest::Client::builder() - .user_agent(subsonic_wire_user_agent()) - .timeout(std::time::Duration::from_secs(120)) - .build() - .map_err(|e| e.to_string())?; - let response = client.get(url).send().await.map_err(|e| e.to_string())?; - if !response.status().is_success() { - return Err(format!("HTTP {}", response.status().as_u16())); - } - let bytes = response.bytes().await.map_err(|e| e.to_string())?; - if bytes.is_empty() { - return Err("empty response".to_string()); - } - enqueue_analysis_seed(app, track_id, &bytes).await -} - -async fn analysis_backfill_worker_loop( - app: tauri::AppHandle, - shared: Arc, - mut wake_rx: tokio::sync::mpsc::UnboundedReceiver<()>, -) { - loop { - if wake_rx.recv().await.is_none() { - break; - } - while let Some((track_id, url)) = { - let mut st = shared - .state - .lock() - .unwrap_or_else(|e| e.into_inner()); - st.try_pop_next() - } { - crate::app_deprintln!("[analysis] backfill worker: start track_id={}", track_id); - let result = analysis_backfill_download_and_seed(&app, &track_id, &url).await; - match &result { - Ok(has_loudness) => crate::app_deprintln!( - "[analysis] backfill ready: {} (has_loudness={})", - track_id, - has_loudness - ), - Err(e) => crate::app_eprintln!("[analysis] backfill failed for {}: {}", track_id, e), - } - let mut st = shared - .state - .lock() - .unwrap_or_else(|e| e.into_inner()); - st.finish_job(&track_id); - } - } -} - -fn analysis_backfill_is_current_track(app: &tauri::AppHandle, track_id: &str) -> bool { - app.try_state::() - .is_some_and(|e| crate::audio::analysis_track_id_is_current_playback(&e, track_id)) -} - -// ─── Full-track waveform + loudness: single CPU worker (mirrors HTTP backfill queue) ─ -// One `spawn_blocking` decode at a time; current playback is high-priority (front + reorder). -// Same `track_id` queued again merges waiters onto one job; while decode runs, same-id -// submitters attach to `running` followers so they all get the same outcome. - -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -enum AnalysisCpuSeedEnqueueKind { - NewBack, - NewFront, - ReorderedFront, - RunningFollower, - MergedQueued, -} - -struct AnalysisCpuSeedJob { - track_id: String, - bytes: Vec, - waiters: Vec>>, -} - -struct AnalysisCpuSeedQueueState { - deque: VecDeque, - /// Decode in progress — same-id callers wait here for the same outcome. - running: Option<( - String, - Arc>>>>, - )>, -} - -impl AnalysisCpuSeedQueueState { - fn enqueue( - &mut self, - track_id: String, - bytes: Vec, - high_priority: bool, - ) -> ( - AnalysisCpuSeedEnqueueKind, - tokio::sync::oneshot::Receiver>, - ) { - let (done_tx, done_rx) = tokio::sync::oneshot::channel(); - let tid = track_id.as_str(); - - if let Some((rtid, followers)) = &self.running { - if rtid == tid { - followers - .lock() - .unwrap_or_else(|e| e.into_inner()) - .push(done_tx); - return (AnalysisCpuSeedEnqueueKind::RunningFollower, done_rx); - } - } - - if let Some(pos) = self.deque.iter().position(|j| j.track_id == track_id) { - let mut job = self.deque.remove(pos).unwrap(); - job.bytes = bytes; - job.waiters.push(done_tx); - let kind = if high_priority { - self.deque.push_front(job); - AnalysisCpuSeedEnqueueKind::ReorderedFront - } else { - self.deque.push_back(job); - AnalysisCpuSeedEnqueueKind::MergedQueued - }; - return (kind, done_rx); - } - - let job = AnalysisCpuSeedJob { - track_id: track_id.clone(), - bytes, - waiters: vec![done_tx], - }; - let kind = if high_priority { - self.deque.push_front(job); - AnalysisCpuSeedEnqueueKind::NewFront - } else { - self.deque.push_back(job); - AnalysisCpuSeedEnqueueKind::NewBack - }; - (kind, done_rx) - } - - fn prune_queued_not_in(&mut self, keep_track_ids: &HashSet<&str>) -> (usize, usize) { - let mut kept = VecDeque::with_capacity(self.deque.len()); - let mut removed_jobs = 0usize; - let mut removed_waiters = 0usize; - while let Some(job) = self.deque.pop_front() { - if keep_track_ids.contains(job.track_id.as_str()) { - kept.push_back(job); - continue; - } - removed_jobs += 1; - removed_waiters += job.waiters.len(); - for tx in job.waiters { - let _ = tx.send(Err( - "cpu-seed pruned: track no longer in playback queue".to_string(), - )); - } - } - self.deque = kept; - (removed_jobs, removed_waiters) - } -} - -struct AnalysisCpuSeedShared { - state: Mutex, - wake_tx: tokio::sync::mpsc::UnboundedSender<()>, -} - -impl Default for AnalysisCpuSeedQueueState { - fn default() -> Self { - Self { - deque: VecDeque::new(), - running: None, - } - } -} - -impl AnalysisCpuSeedShared { - fn ping_worker(&self) { - let _ = self.wake_tx.send(()); - } -} - -static ANALYSIS_CPU_SEED: OnceLock> = OnceLock::new(); - -fn analysis_cpu_seed_shared(app: &tauri::AppHandle) -> Arc { - ANALYSIS_CPU_SEED - .get_or_init(|| { - let (wake_tx, wake_rx) = tokio::sync::mpsc::unbounded_channel(); - let shared = Arc::new(AnalysisCpuSeedShared { - state: Mutex::new(AnalysisCpuSeedQueueState::default()), - wake_tx, - }); - let app = app.clone(); - let sh = shared.clone(); - tauri::async_runtime::spawn(analysis_cpu_seed_worker_loop(app, sh, wake_rx)); - shared - }) - .clone() -} - -/// HTTP backfill + CPU seed queue sizes (debug log only — `app_deprintln!`). -fn emit_analysis_queue_snapshot_line() { - let http = if let Some(arc) = ANALYSIS_BACKFILL.get() { - let st = arc.state.lock().unwrap_or_else(|e| e.into_inner()); - format!( - "http_backfill={{queued:{} download_active:{:?}}}", - st.deque.len(), - st.in_progress.as_deref() - ) - } else { - "http_backfill={{not_started}}".to_string() - }; - - let cpu = if let Some(arc) = ANALYSIS_CPU_SEED.get() { - let st = arc.state.lock().unwrap_or_else(|e| e.into_inner()); - let queued_jobs = st.deque.len(); - let pending_in_queued_jobs: usize = st.deque.iter().map(|j| j.waiters.len()).sum(); - let (decoding_tid, decoding_extra_waiters) = match &st.running { - Some((tid, fl)) => ( - Some(tid.as_str()), - fl.lock().map(|g| g.len()).unwrap_or(0), - ), - None => (None, 0usize), - }; - format!( - "cpu_seed={{queued_jobs:{} pending_channels_in_queue:{} decoding_tid:{:?} extra_waiters_same_id:{}}}", - queued_jobs, - pending_in_queued_jobs, - decoding_tid, - decoding_extra_waiters - ) - } else { - "cpu_seed={{not_started}}".to_string() - }; - - crate::app_deprintln!( - "[analysis] queue_snapshot interval_s=60 note=queues_in_memory_cleared_on_app_restart | {http} | {cpu}" - ); -} - -async fn analysis_queue_snapshot_loop() { - emit_analysis_queue_snapshot_line(); - loop { - tokio::time::sleep(std::time::Duration::from_secs(60)).await; - emit_analysis_queue_snapshot_line(); - } -} - -async fn analysis_cpu_seed_worker_loop( - app: tauri::AppHandle, - shared: Arc, - mut wake_rx: tokio::sync::mpsc::UnboundedReceiver<()>, -) { - loop { - if wake_rx.recv().await.is_none() { - break; - } - loop { - let (job, followers) = { - let mut st = shared.state.lock().unwrap_or_else(|e| e.into_inner()); - let Some(j) = st.deque.pop_front() else { - break; - }; - let fl = Arc::new(Mutex::new(Vec::new())); - st.running = Some((j.track_id.clone(), fl.clone())); - (j, fl) - }; - let tid_log = job.track_id.clone(); - let app2 = app.clone(); - let tid = job.track_id.clone(); - let bytes = job.bytes; - let outcome = tokio::task::spawn_blocking(move || { - analysis_cache::seed_from_bytes_execute(&app2, &tid, &bytes) - }) - .await - .unwrap_or_else(|e| Err(format!("cpu-seed spawn_blocking: {e}"))); - - let mut extra = followers - .lock() - .unwrap_or_else(|e| e.into_inner()) - .drain(..) - .collect::>(); - for tx in job.waiters { - let _ = tx.send(outcome.clone()); - } - for tx in extra.drain(..) { - let _ = tx.send(outcome.clone()); - } - - { - let mut st = shared.state.lock().unwrap_or_else(|e| e.into_inner()); - st.running = None; - } - let ok = outcome.as_ref().map(|o| *o == analysis_cache::SeedFromBytesOutcome::Upserted).unwrap_or(false); - crate::app_deprintln!( - "[analysis] cpu-seed worker: done track_id={} upserted={}", - tid_log, - ok - ); - } - } -} - -/// Submit full-buffer analysis; serializes with other producers. `high_priority` mirrors -/// HTTP backfill head insertion for the currently playing track. -/// -/// Emits `analysis:waveform-updated` when analysis **wrote** new waveform data (`Upserted`). -/// Cache-hit skips (`SkippedWaveformCacheHit`) omit the event so the frontend does not -/// re-run loudness refresh / waveform IPC for rows that were already current. -pub(crate) async fn submit_analysis_cpu_seed( - app: tauri::AppHandle, - track_id: String, - bytes: Vec, - high_priority: bool, -) -> Result { - let shared = analysis_cpu_seed_shared(&app); - let rx = { - let mut st = shared.state.lock().unwrap_or_else(|e| e.into_inner()); - let (kind, rx) = st.enqueue(track_id.clone(), bytes, high_priority); - crate::app_deprintln!("[analysis] cpu-seed submit: kind={kind:?} high_priority={high_priority}"); - drop(st); - shared.ping_worker(); - rx - }; - let outcome = match rx.await { - Ok(res) => res?, - Err(_) => return Err("cpu-seed: result channel dropped".to_string()), - }; - if matches!(outcome, analysis_cache::SeedFromBytesOutcome::Upserted) { - let _ = app.emit( - "analysis:waveform-updated", - WaveformUpdatedPayload { - track_id: track_id.clone(), - is_partial: false, - }, - ); - } - Ok(outcome) -} - /// Holds the live system-tray icon handle. `None` means the tray is currently hidden/removed. /// Dropping the inner `TrayIcon` fully removes it from the OS notification area on all platforms. type TrayState = Mutex>; @@ -588,13 +145,6 @@ struct WaveformCachePayload { updated_at: i64, } -#[derive(Clone, serde::Serialize)] -#[serde(rename_all = "camelCase")] -struct WaveformUpdatedPayload { - track_id: String, - is_partial: bool, -} - #[derive(serde::Serialize)] #[serde(rename_all = "camelCase")] struct LoudnessCachePayload { diff --git a/src-tauri/src/lib_commands/app_api/analysis.rs b/src-tauri/src/lib_commands/app_api/analysis.rs index 54d9f5a3..525fbeca 100644 --- a/src-tauri/src/lib_commands/app_api/analysis.rs +++ b/src-tauri/src/lib_commands/app_api/analysis.rs @@ -205,25 +205,8 @@ pub(crate) fn analysis_prune_pending_to_track_ids( } let keep_track_ids: HashSet<&str> = normalized.iter().map(|s| s.as_str()).collect(); - let http_removed = if let Some(shared) = ANALYSIS_BACKFILL.get() { - let mut st = shared - .state - .lock() - .map_err(|_| "analysis backfill lock poisoned".to_string())?; - st.prune_queued_not_in(&keep_track_ids) - } else { - 0 - }; - - let (cpu_removed_jobs, cpu_removed_waiters) = if let Some(shared) = ANALYSIS_CPU_SEED.get() { - let mut st = shared - .state - .lock() - .map_err(|_| "analysis cpu-seed lock poisoned".to_string())?; - st.prune_queued_not_in(&keep_track_ids) - } else { - (0, 0) - }; + let (http_removed, cpu_removed_jobs, cpu_removed_waiters) = + prune_analysis_queues(&keep_track_ids)?; if http_removed > 0 || cpu_removed_jobs > 0 { crate::app_deprintln!(