use std::path::{Path, PathBuf}; use doctate_common::oneliners::OnelinerState; use tracing::{info, warn}; use super::{TranscribeJob, TranscribeSender}; use crate::config::Config; use crate::events::EventSender; use crate::gazetteer::Gazetteer; use crate::paths; /// Walk `data_path/*/` and enqueue every recording pending transcription for /// every user. Intended for startup; the same primitive backs the per-user /// self-heal triggered by web handlers. pub async fn scan_and_enqueue(data_path: &Path, tx: &TranscribeSender) { let Ok(mut users) = tokio::fs::read_dir(data_path).await else { return; }; // Collect + sort by slug so enqueue order is deterministic across // users (within a user, the primitive sorts by timestamp). let mut user_dirs: Vec<(String, std::path::PathBuf)> = Vec::new(); while let Ok(Some(user_entry)) = users.next_entry().await { if let Some(slug) = user_entry.file_name().to_str().map(str::to_owned) { user_dirs.push((slug, user_entry.path())); } } user_dirs.sort_by(|a, b| a.0.cmp(&b.0)); let mut total = 0; for (slug, path) in user_dirs { total += enqueue_pending_for_user(&path, &slug, tx).await; } info!(count = total, "Transcribe recovery scan finished"); } /// Scan a single `//*.m4a` set and enqueue every recording /// whose `.transcript.txt` is missing. Returns the number sent. Channel send /// is awaited (bounded channel); only fails if the worker is gone. pub async fn enqueue_pending_for_user( user_root: &Path, slug: &str, tx: &TranscribeSender, ) -> usize { let Ok(mut cases) = tokio::fs::read_dir(user_root).await else { return 0; }; let mut pending: Vec = Vec::new(); while let Ok(Some(case_entry)) = cases.next_entry().await { let case_dir = case_entry.path(); if !case_dir.is_dir() || crate::paths::is_deleted(&case_dir).await { continue; } let Ok(mut files) = tokio::fs::read_dir(&case_dir).await else { continue; }; while let Ok(Some(file)) = files.next_entry().await { let path = file.path(); if path.extension().and_then(|s| s.to_str()) != Some("m4a") { continue; } if path.with_extension("transcript.txt").exists() { continue; } pending.push(path); } } // Oldest first — filenames are ISO timestamps so path order = chronological. pending.sort(); let mut sent = 0; for path in pending { if tx .send(TranscribeJob { audio_path: path.clone(), user_slug: slug.to_owned(), }) .await .is_err() { warn!(file = %path.display(), "recovery send failed (worker gone)"); return sent; } sent += 1; } sent } /// Walk `data_path/*/*` and return every case directory that has at least /// one non-empty `*.transcript.txt` and whose persisted oneliner state /// warrants (re-)generation, paired with the user-slug (parent directory /// name). Pure filesystem scan, no LLM calls — separated from the /// regeneration wrapper so it can be unit-tested without mocking Ollama. /// /// Retry policy by persisted state: /// - missing file → retry (first generation, or crash before write) /// - [`OnelinerState::Error`] → retry (transient LLM failure last time) /// - [`OnelinerState::Ready`] or [`OnelinerState::Empty`] → do **not** /// retry; `Empty` is a valid outcome and `Ready` is already correct. /// A fresh transcript write will trigger regeneration through the /// normal worker path. pub(crate) async fn cases_needing_oneliner(data_path: &Path) -> Vec<(PathBuf, String)> { let mut out: Vec<(PathBuf, String)> = Vec::new(); let Ok(mut users) = tokio::fs::read_dir(data_path).await else { return out; }; while let Ok(Some(user_entry)) = users.next_entry().await { let user_path = user_entry.path(); if !user_path.is_dir() { continue; } let Some(slug) = user_entry.file_name().to_str().map(str::to_owned) else { continue; }; for case_dir in cases_needing_oneliner_in(&user_path).await { out.push((case_dir, slug.clone())); } } out.sort_by(|a, b| a.0.cmp(&b.0)); out } /// Per-user variant of [`cases_needing_oneliner`]: scans only `user_root`'s /// case dirs. Used both as the building block for the all-users scan and /// from the page-load self-heal path. pub(crate) async fn cases_needing_oneliner_in(user_root: &Path) -> Vec { let mut out: Vec = Vec::new(); let Ok(mut cases) = tokio::fs::read_dir(user_root).await else { return out; }; while let Ok(Some(case_entry)) = cases.next_entry().await { let case_dir = case_entry.path(); if !case_dir.is_dir() || paths::is_deleted(&case_dir).await { continue; } let (state, _) = paths::read_oneliner_state(&case_dir).await; let needs_retry = match state { None | Some(OnelinerState::Error { .. }) => true, Some(OnelinerState::Ready { .. } | OnelinerState::Empty { .. }) => false, }; if !needs_retry { continue; } if has_non_empty_transcript(&case_dir).await { out.push(case_dir); } } out.sort(); out } async fn has_non_empty_transcript(case_dir: &Path) -> bool { let Ok(mut files) = tokio::fs::read_dir(case_dir).await else { return false; }; while let Ok(Some(f)) = files.next_entry().await { let path = f.path(); let Some(name) = path.file_name().and_then(|s| s.to_str()) else { continue; }; if !name.ends_with(".transcript.txt") { continue; } if let Ok(text) = tokio::fs::read_to_string(&path).await && !text.trim().is_empty() { return true; } } false } /// Regenerate `oneliner.json` for every case that has transcripts but no /// state file (crash between transcript-write and oneliner-write) or a /// previously-errored state. Called at startup after `scan_and_enqueue`. /// Sequential on purpose — the oneliner model is light, but spamming it /// in parallel is not worth it. pub async fn regenerate_missing_oneliners( data_path: &Path, client: &reqwest::Client, config: &Config, vocab: &Gazetteer, events_tx: &EventSender, ) { let cases = cases_needing_oneliner(data_path).await; if cases.is_empty() { return; } info!(count = cases.len(), "Regenerating missing oneliners"); for (case_dir, slug) in cases { super::worker::update_oneliner(&case_dir, &slug, client, config, vocab, events_tx).await; } } /// Per-user variant of [`regenerate_missing_oneliners`]: fixes only the /// given user's cases. Intended for the page-load self-heal path, so a /// single missing oneliner is restored as soon as its owner opens a page /// — no server restart needed. pub async fn regenerate_missing_oneliners_for_user( user_root: &Path, slug: &str, client: &reqwest::Client, config: &Config, vocab: &Gazetteer, events_tx: &EventSender, ) { let cases = cases_needing_oneliner_in(user_root).await; if cases.is_empty() { return; } for case_dir in cases { super::worker::update_oneliner(&case_dir, slug, client, config, vocab, events_tx).await; } } #[cfg(test)] mod tests { use super::*; use tempfile::tempdir; #[tokio::test] async fn cases_needing_oneliner_finds_crash_survivor() { let data = tempdir().unwrap(); let case = data.path().join("user").join("case1"); tokio::fs::create_dir_all(&case).await.unwrap(); tokio::fs::write(case.join("2026-04-16T10-00-00Z.transcript.txt"), "content") .await .unwrap(); let got = cases_needing_oneliner(data.path()).await; assert_eq!(got, vec![(case, "user".to_owned())]); } async fn write_state(case: &Path, state: &OnelinerState) { let bytes = serde_json::to_vec(state).unwrap(); tokio::fs::write(case.join("oneliner.json"), bytes) .await .unwrap(); } #[tokio::test] async fn cases_needing_oneliner_skips_ready_state() { let data = tempdir().unwrap(); let case = data.path().join("user").join("case1"); tokio::fs::create_dir_all(&case).await.unwrap(); tokio::fs::write(case.join("2026-04-16T10-00-00Z.transcript.txt"), "content") .await .unwrap(); write_state( &case, &OnelinerState::Ready { text: "Existing".into(), generated_at: "2026-04-16T10:01:00Z".into(), }, ) .await; assert!(cases_needing_oneliner(data.path()).await.is_empty()); } #[tokio::test] async fn cases_needing_oneliner_skips_empty_state() { let data = tempdir().unwrap(); let case = data.path().join("user").join("case1"); tokio::fs::create_dir_all(&case).await.unwrap(); tokio::fs::write(case.join("2026-04-16T10-00-00Z.transcript.txt"), "content") .await .unwrap(); write_state( &case, &OnelinerState::Empty { generated_at: "2026-04-16T10:01:00Z".into(), }, ) .await; assert!(cases_needing_oneliner(data.path()).await.is_empty()); } #[tokio::test] async fn cases_needing_oneliner_retries_error_state() { let data = tempdir().unwrap(); let case = data.path().join("user").join("case1"); tokio::fs::create_dir_all(&case).await.unwrap(); tokio::fs::write(case.join("2026-04-16T10-00-00Z.transcript.txt"), "content") .await .unwrap(); write_state( &case, &OnelinerState::Error { generated_at: "2026-04-16T10:01:00Z".into(), }, ) .await; let got = cases_needing_oneliner(data.path()).await; assert_eq!(got, vec![(case, "user".to_owned())]); } #[tokio::test] async fn cases_needing_oneliner_skips_cases_with_only_empty_transcripts() { let data = tempdir().unwrap(); let case = data.path().join("user").join("case1"); tokio::fs::create_dir_all(&case).await.unwrap(); tokio::fs::write(case.join("2026-04-16T10-00-00Z.transcript.txt"), " ") .await .unwrap(); assert!(cases_needing_oneliner(data.path()).await.is_empty()); } #[tokio::test] async fn cases_needing_oneliner_skips_deleted_cases() { let data = tempdir().unwrap(); let case = data.path().join("user").join("case1"); tokio::fs::create_dir_all(&case).await.unwrap(); tokio::fs::write(case.join("2026-04-16T10-00-00Z.transcript.txt"), "content") .await .unwrap(); // Empty JSON object is a valid delete marker (see paths::is_deleted). tokio::fs::write(case.join(".deleted"), "{}").await.unwrap(); assert!(cases_needing_oneliner(data.path()).await.is_empty()); } }