fix(analysis): persist failed tracks and reconcile progress counts (#867)

* fix(analysis): persist failed-track suppression and reduce aggressive polling

Persist unsupported/broken analysis tracks as failed entries and expose them in Settings with track metadata, export, and targeted rescan actions. Also make aggressive-mode completion checks cheap by gating re-entry on live track-count changes with a startup seed and 5-minute recheck cadence.

* fix(analysis): mark unsupported decode tracks as failed in cpu-seed

When full-seed falls back to waveform-only (no EBU loudness) or enrichment decode fails, persist analysis_track status as failed so aggressive backfill does not requeue the same unsupported tracks indefinitely.

* fix(analysis): reconcile legacy ready tracks without loudness

Auto-mark legacy ready tracks that only miss loudness as failed during needs-work checks so analysis progress converges instead of staying permanently pending.

* docs(changelog): add failed-analysis recovery notes for PR 867

Document persistent failed-track handling, analytics strategy controls, and low-cost aggressive-mode recheck behavior in the 1.47.0 changelog.

* fix(i18n): align settings locale coverage across all languages

Sync missing Settings keys for all shipped locales and replace recent English fallbacks with localized strings so analytics and backup UI text stays consistent outside en/ru.
This commit is contained in:
cucadmuh
2026-05-25 01:37:15 +03:00
committed by GitHub
parent 820f71c421
commit bc85065316
21 changed files with 1263 additions and 17 deletions
@@ -11,6 +11,7 @@ use symphonia::core::meta::MetadataOptions;
use symphonia::core::probe::Hint;
use symphonia::core::units::Time;
use tauri::{Manager, Runtime};
use psysonic_core::track_enrichment::TrackEnrichmentOutcome;
use crate::analysis_perf::AnalysisSeedTimings;
@@ -78,12 +79,32 @@ pub fn seed_from_bytes_execute<R: Runtime>(
}
let bpm_ms = if !server_id.is_empty() {
let bpm_started = Instant::now();
let _ = crate::track_enrichment::run_track_enrichment_if_needed(
let enrichment_outcome = crate::track_enrichment::run_track_enrichment_if_needed(
app,
server_id,
track_id,
bytes,
);
if matches!(enrichment_outcome, TrackEnrichmentOutcome::Failed) {
let key = TrackKey {
server_id: server_id.to_string(),
track_id: track_id.to_string(),
md5_16kb: md5_16kb.clone(),
};
let _ = cache.touch_track_status(&key, "failed");
}
if matches!(outcome, SeedFromBytesOutcome::Upserted) {
if let Ok(coverage) = cache.content_cache_coverage(server_id, track_id, &md5_16kb) {
if !coverage.has_loudness {
let key = TrackKey {
server_id: server_id.to_string(),
track_id: track_id.to_string(),
md5_16kb: md5_16kb.clone(),
};
let _ = cache.touch_track_status(&key, "failed");
}
}
}
bpm_started.elapsed().as_millis() as u64
} else {
0
@@ -199,6 +220,7 @@ pub fn seed_from_bytes_into_cache(
);
}
Err(e) => {
let _ = cache.touch_track_status(&key, "failed");
crate::app_deprintln!(
"[analysis] full-track analysis failed track_id={} elapsed_ms={} err={}",
track_id,
@@ -7,5 +7,6 @@ pub use compute::{
seed_from_bytes_execute, seed_from_bytes_into_cache, PcmAnalysisWindow, SeedFromBytesOutcome,
};
pub use store::{
AnalysisCache, AnalysisDeleteServerReport, LoudnessEntry, TrackKey, WaveformEntry,
AnalysisCache, AnalysisDeleteServerReport, FailedTrackEntry, LoudnessEntry, TrackKey,
WaveformEntry,
};
@@ -92,6 +92,13 @@ pub struct AnalysisDeleteServerReport {
pub loudness: u64,
}
#[derive(Debug, Clone)]
pub struct FailedTrackEntry {
pub track_id: String,
pub md5_16kb: String,
pub updated_at: i64,
}
pub struct AnalysisCache {
conn: Mutex<Connection>,
}
@@ -110,6 +117,10 @@ fn track_id_cache_variants(id: &str) -> Vec<String> {
out
}
fn normalize_track_id(id: &str) -> String {
id.strip_prefix("stream:").unwrap_or(id).to_string()
}
pub(super) fn now_unix_ts() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
@@ -538,6 +549,66 @@ impl AnalysisCache {
query_latest_md5_16kb_scoped(&conn, server_id, track_id)
}
/// Latest analysis status row for `(server_id, track_id)` (tries bare and
/// `stream:` variants). Used to suppress infinite retries for tracks that
/// repeatedly fail decode/enrichment.
pub fn get_latest_status_for_track(
&self,
server_id: &str,
track_id: &str,
) -> Result<Option<(String, i64)>, String> {
let conn = self.conn.lock().map_err(|_| "analysis_cache lock poisoned".to_string())?;
query_latest_status_scoped(&conn, server_id, track_id)
}
pub fn count_failed_tracks(&self, server_id: &str) -> Result<i64, String> {
let conn = self.conn.lock().map_err(|_| "analysis_cache lock poisoned".to_string())?;
query_failed_tracks_count_scoped(&conn, server_id)
}
pub fn list_failed_tracks(
&self,
server_id: &str,
limit: Option<usize>,
) -> Result<Vec<FailedTrackEntry>, String> {
let conn = self.conn.lock().map_err(|_| "analysis_cache lock poisoned".to_string())?;
query_failed_tracks_scoped(&conn, server_id, limit)
}
pub fn clear_failed_tracks(
&self,
server_id: &str,
track_ids: &[String],
) -> Result<u64, String> {
let conn = self.conn.lock().map_err(|_| "analysis_cache lock poisoned".to_string())?;
if track_ids.is_empty() {
let deleted = conn
.execute(
"DELETE FROM analysis_track WHERE server_id = ?1 AND status = 'failed'",
params![server_id],
)
.map_err(|e| e.to_string())?;
return Ok(deleted as u64);
}
let mut total = 0u64;
for id in track_ids {
let tid = id.trim();
if tid.is_empty() {
continue;
}
for variant in track_id_cache_variants(tid) {
let deleted = conn
.execute(
"DELETE FROM analysis_track WHERE server_id = ?1 AND track_id = ?2 AND status = 'failed'",
params![server_id, variant],
)
.map_err(|e| e.to_string())?;
total = total.saturating_add(deleted as u64);
}
}
Ok(total)
}
/// Both waveform and loudness rows exist for this `(server_id, track_id)` —
/// a CPU seed from bytes/file would only decode the file to immediately skip
/// with `SkippedWaveformCacheHit`.
@@ -637,6 +708,104 @@ fn query_latest_md5_16kb_scoped(
Ok(None)
}
fn query_latest_status_scoped(
conn: &Connection,
server_id: &str,
track_id: &str,
) -> Result<Option<(String, i64)>, String> {
const SQL: &str = r#"
SELECT status, updated_at
FROM analysis_track
WHERE server_id = ?1
AND track_id = ?2
ORDER BY updated_at DESC
LIMIT 1
"#;
let mut latest: Option<(String, i64)> = None;
for tid in track_id_cache_variants(track_id) {
let row = conn
.query_row(SQL, params![server_id, tid], |row| {
Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?))
})
.optional()
.map_err(|e| e.to_string())?;
if let Some(candidate) = row {
let take = latest
.as_ref()
.is_none_or(|(_, ts)| candidate.1 > *ts);
if take {
latest = Some(candidate);
}
}
}
Ok(latest)
}
fn query_failed_tracks_count_scoped(conn: &Connection, server_id: &str) -> Result<i64, String> {
conn.query_row(
r#"
SELECT COUNT(DISTINCT normalized_track_id)
FROM (
SELECT CASE
WHEN track_id LIKE 'stream:%' THEN SUBSTR(track_id, 8)
ELSE track_id
END AS normalized_track_id
FROM analysis_track
WHERE server_id = ?1
AND status = 'failed'
)
WHERE normalized_track_id IS NOT NULL
AND normalized_track_id != ''
"#,
params![server_id],
|row| row.get(0),
)
.map_err(|e| e.to_string())
}
fn query_failed_tracks_scoped(
conn: &Connection,
server_id: &str,
limit: Option<usize>,
) -> Result<Vec<FailedTrackEntry>, String> {
const SQL: &str = r#"
SELECT track_id, md5_16kb, updated_at
FROM analysis_track
WHERE server_id = ?1
AND status = 'failed'
ORDER BY updated_at DESC
"#;
let mut stmt = conn.prepare(SQL).map_err(|e| e.to_string())?;
let rows = stmt
.query_map(params![server_id], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, i64>(2)?,
))
})
.map_err(|e| e.to_string())?;
let mut out: Vec<FailedTrackEntry> = Vec::new();
let mut seen = std::collections::HashSet::<String>::new();
for row in rows {
let (track_id_raw, md5_16kb, updated_at) = row.map_err(|e| e.to_string())?;
let normalized = normalize_track_id(track_id_raw.trim());
if normalized.is_empty() || !seen.insert(normalized.clone()) {
continue;
}
out.push(FailedTrackEntry {
track_id: normalized,
md5_16kb,
updated_at,
});
if limit.is_some_and(|n| out.len() >= n) {
break;
}
}
Ok(out)
}
/// Server-scoped variant of the "latest loudness for this track" lookup.
fn query_latest_loudness_scoped(
conn: &Connection,
@@ -1673,4 +1842,62 @@ mod tests {
let cache = AnalysisCache::init(&handle).expect("analysis cache init with mock app");
cache.checkpoint_wal("init-test").unwrap();
}
#[test]
fn latest_status_for_track_prefers_newest_variant_timestamp() {
let cache = AnalysisCache::open_in_memory();
let base = key_on("server-a", "track-1");
let prefixed = key_on("server-a", "stream:track-1");
cache.touch_track_status(&base, "queued").unwrap();
std::thread::sleep(std::time::Duration::from_millis(1200));
cache.touch_track_status(&prefixed, "failed").unwrap();
let row = cache
.get_latest_status_for_track("server-a", "track-1")
.unwrap()
.expect("latest status row");
assert_eq!(row.0, "failed");
}
#[test]
fn failed_track_queries_deduplicate_stream_variants() {
let cache = AnalysisCache::open_in_memory();
let base = key_on("server-a", "track-2");
let prefixed = key_on("server-a", "stream:track-2");
let other = key_on("server-a", "track-3");
cache.touch_track_status(&base, "failed").unwrap();
std::thread::sleep(std::time::Duration::from_millis(1200));
cache.touch_track_status(&prefixed, "failed").unwrap();
cache.touch_track_status(&other, "failed").unwrap();
let count = cache.count_failed_tracks("server-a").unwrap();
assert_eq!(count, 2);
let listed = cache.list_failed_tracks("server-a", None).unwrap();
assert_eq!(listed.len(), 2);
assert!(listed.iter().any(|r| r.track_id == "track-2"));
assert!(listed.iter().any(|r| r.track_id == "track-3"));
}
#[test]
fn clear_failed_tracks_removes_only_failed_rows() {
let cache = AnalysisCache::open_in_memory();
let failed = key_on("server-a", "track-failed");
let ready = key_on("server-a", "track-ready");
cache.touch_track_status(&failed, "failed").unwrap();
cache.touch_track_status(&ready, "ready").unwrap();
let deleted = cache
.clear_failed_tracks("server-a", &["track-failed".to_string()])
.unwrap();
assert_eq!(deleted, 1);
assert_eq!(cache.count_failed_tracks("server-a").unwrap(), 0);
let ready_latest = cache
.get_latest_status_for_track("server-a", "track-ready")
.unwrap()
.expect("ready row stays");
assert_eq!(ready_latest.0, "ready");
}
}
@@ -5,6 +5,7 @@ use std::sync::{Arc, Mutex, OnceLock};
use tauri::{Emitter, Manager};
use psysonic_core::user_agent::subsonic_wire_user_agent;
use psysonic_core::track_enrichment::TrackEnrichmentOutcome;
use crate::analysis_cache;
@@ -435,7 +436,18 @@ async fn enqueue_track_analysis_with_fetch(
content_hash
);
let bpm_started = std::time::Instant::now();
run_track_enrichment_from_bytes(app, server_id, track_id, bytes).await;
let outcome = run_track_enrichment_from_bytes(app, server_id, track_id, bytes).await;
if matches!(outcome, TrackEnrichmentOutcome::Failed) {
if let Some(cache) = app.try_state::<analysis_cache::AnalysisCache>() {
let key = analysis_cache::TrackKey {
server_id: server_id.to_string(),
track_id: track_id.to_string(),
md5_16kb: content_hash.clone(),
};
let _ = cache.touch_track_status(&key, "failed");
}
return Err("track enrichment failed".to_string());
}
let bpm_ms = bpm_started.elapsed().as_millis() as u64;
emit_analysis_track_perf(app, track_id, fetch_ms, 0, bpm_ms);
return Ok(EnqueueTrackAnalysisOutcome::RanEnrichmentOnly);
@@ -452,18 +464,22 @@ pub async fn run_track_enrichment_from_bytes(
server_id: &str,
track_id: &str,
bytes: &[u8],
) {
) -> TrackEnrichmentOutcome {
if server_id.is_empty() {
return;
return TrackEnrichmentOutcome::SkippedNoServer;
}
let app = app.clone();
let sid = server_id.to_string();
let tid = track_id.to_string();
let data = bytes.to_vec();
let _ = tokio::task::spawn_blocking(move || {
crate::track_enrichment::run_track_enrichment_if_needed(&app, &sid, &tid, &data);
match tokio::task::spawn_blocking(move || {
crate::track_enrichment::run_track_enrichment_if_needed(&app, &sid, &tid, &data)
})
.await;
.await
{
Ok(outcome) => outcome,
Err(_) => TrackEnrichmentOutcome::Failed,
}
}
/// Read a local file and run [`enqueue_track_analysis`] (hot cache, offline, spill promote).
@@ -59,6 +59,14 @@ pub struct AnalysisDeleteServerReportDto {
pub loudness: u64,
}
#[derive(Debug, Clone, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct AnalysisFailedTrackDto {
pub track_id: String,
pub md5_16kb: String,
pub updated_at: i64,
}
impl From<analysis_cache::AnalysisDeleteServerReport> for AnalysisDeleteServerReportDto {
fn from(value: analysis_cache::AnalysisDeleteServerReport) -> Self {
Self {
@@ -69,6 +77,16 @@ impl From<analysis_cache::AnalysisDeleteServerReport> for AnalysisDeleteServerRe
}
}
impl From<analysis_cache::FailedTrackEntry> for AnalysisFailedTrackDto {
fn from(value: analysis_cache::FailedTrackEntry) -> Self {
Self {
track_id: value.track_id,
md5_16kb: value.md5_16kb,
updated_at: value.updated_at,
}
}
}
#[derive(serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AnalysisServerKeyMigrationDto {
@@ -227,6 +245,54 @@ pub fn analysis_delete_all_for_server(
Ok(report.into())
}
#[tauri::command]
pub fn analysis_get_failed_track_count(
server_id: String,
cache: tauri::State<'_, analysis_cache::AnalysisCache>,
) -> Result<i64, String> {
let server_id = server_id.trim().to_string();
if server_id.is_empty() {
return Ok(0);
}
cache.count_failed_tracks(&server_id)
}
#[tauri::command]
pub fn analysis_list_failed_tracks(
server_id: String,
limit: Option<u32>,
cache: tauri::State<'_, analysis_cache::AnalysisCache>,
) -> Result<Vec<AnalysisFailedTrackDto>, String> {
let server_id = server_id.trim().to_string();
if server_id.is_empty() {
return Ok(Vec::new());
}
let limit = limit
.map(|v| usize::try_from(v).unwrap_or(usize::MAX))
.map(|v| v.clamp(1, 5_000));
let rows = cache.list_failed_tracks(&server_id, limit)?;
Ok(rows.into_iter().map(AnalysisFailedTrackDto::from).collect())
}
#[tauri::command]
pub fn analysis_clear_failed_tracks(
server_id: String,
track_ids: Option<Vec<String>>,
cache: tauri::State<'_, analysis_cache::AnalysisCache>,
) -> Result<u64, String> {
let server_id = server_id.trim().to_string();
if server_id.is_empty() {
return Err("server_id required".to_string());
}
let track_ids = track_ids
.unwrap_or_default()
.into_iter()
.map(|id| id.trim().to_string())
.filter(|id| !id.is_empty())
.collect::<Vec<_>>();
cache.clear_failed_tracks(&server_id, &track_ids)
}
#[tauri::command]
pub fn analysis_migrate_server_index_keys(
mappings: Vec<AnalysisServerKeyMigrationDto>,
@@ -7,7 +7,7 @@ use psysonic_core::track_analysis::TrackAnalysisPlan;
use psysonic_core::track_enrichment::TrackEnrichmentPort;
use tauri::{AppHandle, Manager};
use crate::analysis_cache::AnalysisCache;
use crate::analysis_cache::{AnalysisCache, TrackKey};
pub fn plan_track_analysis(
app: &AppHandle,
@@ -52,6 +52,40 @@ pub fn track_analysis_needs_work(
server_id: &str,
track_id: &str,
) -> Result<bool, String> {
if let Some(cache) = app.try_state::<AnalysisCache>() {
let latest_status = cache.get_latest_status_for_track(server_id, track_id)?;
if latest_status
.as_ref()
.is_some_and(|(status, _)| status == "failed")
{
return Ok(false);
}
let plan = plan_track_analysis_from_cache(app, server_id, track_id)?;
if !plan.any() {
return Ok(false);
}
// Legacy reconciliation: some old rows are persisted as `ready` with
// waveform present but no loudness (typically unsupported decode path).
// Those tracks spin forever in pending without converging. Promote to
// terminal `failed` so scheduler/progress can converge.
if latest_status
.as_ref()
.is_some_and(|(status, _)| status == "ready")
&& plan.need_loudness
&& !plan.need_waveform
{
if let Some(md5) = cache.get_latest_md5_16kb_for_track(server_id, track_id)? {
let key = TrackKey {
server_id: server_id.to_string(),
track_id: track_id.to_string(),
md5_16kb: md5,
};
let _ = cache.touch_track_status(&key, "failed");
}
return Ok(false);
}
return Ok(plan.any());
}
Ok(plan_track_analysis_from_cache(app, server_id, track_id)?.any())
}
@@ -159,4 +193,5 @@ mod tests {
let (wf, ld) = cache_gaps_for_content(Some(&cache), "s1", "t1", "abc");
assert!(!wf && !ld, "bare id should resolve stream: cached fingerprint");
}
}