//! End-to-end regression test for the transient-transcribe-failure heal. //! //! Scenario: Minerva is briefly offline (HTTP 503) when the doctor's //! recording arrives. The transcribe worker classifies the error as //! transient and leaves the audio as plain `.m4a` — no `.failed` //! suffix. On the next page-load, `heal_orphans_if_idle` re-enqueues //! the pending `.m4a`. By then Minerva is back; Whisper returns the //! transcript and the case advances normally. //! //! Pre-fix behaviour: any Whisper error renamed `.m4a` → `.m4a.failed`, //! terminating progress permanently. Only a manual admin reset un-failed //! the recording. use std::path::PathBuf; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use std::time::Duration; use doctate_server::config::{Config, User}; use doctate_server::events; use doctate_server::gazetteer::Gazetteer; use doctate_server::{ AnalyzeBusy, OnelinerHealBusy, PipelineState, TranscribeBusy, WorkerBusy, analyze, transcribe, }; use tempfile::tempdir; use wiremock::matchers::{method, path as wm_path}; use wiremock::{Mock, MockServer, ResponseTemplate}; fn fixture(name: &str) -> PathBuf { PathBuf::from(env!("CARGO_MANIFEST_DIR")) .join("tests/fixtures") .join(name) } #[tokio::test] async fn transient_whisper_failure_is_reenqueued_and_recovers() { // Arrange filesystem: one real m4a for user "dr_test". let data = tempdir().unwrap(); let slug = "dr_test"; let user_root = data.path().join(slug); let case_dir = user_root.join("case-transient"); std::fs::create_dir_all(&case_dir).unwrap(); let audio = case_dir.join("2026-04-15T20-30-00Z.m4a"); std::fs::copy(fixture("sample.m4a"), &audio).unwrap(); // Arrange Whisper mock: first call → 503 (Minerva down), // subsequent calls → 200 with a transcript. wiremock matches // mocks in registration order; `up_to_n_times(1)` retires the // first mock after one hit, so the second takes over for any // retry. let whisper = MockServer::start().await; Mock::given(method("POST")) .and(wm_path("/asr")) .respond_with(ResponseTemplate::new(503).set_body_string("service unavailable")) .up_to_n_times(1) .mount(&whisper) .await; Mock::given(method("POST")) .and(wm_path("/asr")) .respond_with(ResponseTemplate::new(200).set_body_string("recovered text")) .mount(&whisper) .await; // Config points at the mock Whisper, ollama intentionally unreachable // (the oneliner heal may fire in parallel; we only assert on // the transcript artefact, not on the oneliner). let mut cfg = Config::test_default(); cfg.data_path = data.path().to_path_buf(); cfg.whisper_url = whisper.uri(); cfg.whisper_timeout_seconds = 5; cfg.users = vec![User { slug: slug.into(), api_key: "k".into(), web_password: "unused".into(), role: "doctor".into(), whisper: Default::default(), retention: Default::default(), window_hours: 72, preview_lines: 2, }]; let config = Arc::new(cfg); let vocab = Arc::new(Gazetteer::empty()); let events_tx = events::channel(); let http_client = reqwest::Client::new(); let (tx_a, _rx_a) = analyze::channel(); let (tx_t, rx_t) = transcribe::channel(); let transcribe_busy: WorkerBusy = Arc::new(AtomicBool::new(false)); let heal_busy: WorkerBusy = Arc::new(AtomicBool::new(false)); // Spawn the transcribe worker in the background; it processes jobs // until the channel is closed at test teardown. let worker = tokio::spawn(transcribe::worker::run( rx_t, config.clone(), http_client.clone(), transcribe_busy.clone(), vocab.clone(), events_tx.clone(), )); let pipeline = PipelineState { analyze_busy: AnalyzeBusy(Arc::new(AtomicBool::new(false))), analyze_tx: tx_a, transcribe_busy: TranscribeBusy(transcribe_busy.clone()), transcribe_tx: tx_t.clone(), oneliner_heal_busy: OnelinerHealBusy(heal_busy.clone()), }; // Act 1: enqueue the initial job, wait for the 503 failure. tx_t.send(transcribe::TranscribeJob { audio_path: audio.clone(), user_slug: slug.into(), }) .await .unwrap(); // Wait until the worker has processed the first job (busy false // AND the 503 mock has recorded a hit). Polling both conditions // avoids a race where the worker hasn't started yet. let start = std::time::Instant::now(); loop { let idle = !transcribe_busy.load(Ordering::Acquire); let hits = whisper.received_requests().await.unwrap().len(); if idle && hits >= 1 { break; } if start.elapsed() > Duration::from_secs(10) { panic!("worker did not process first job within 10s (idle={idle}, hits={hits})"); } tokio::time::sleep(Duration::from_millis(25)).await; } // Invariant after transient failure: .m4a still there, no .failed. assert!(audio.exists(), "transient error must not consume the .m4a"); assert!( !case_dir.join("2026-04-15T20-30-00Z.m4a.failed").exists(), ".m4a.failed must not exist after transient error" ); assert!( !case_dir .join("2026-04-15T20-30-00Z.transcript.txt") .exists(), "no transcript yet — the 503 produced nothing" ); // Act 2: heal re-enqueues the pending .m4a (second mock → 200). pipeline .heal_orphans_if_idle(&user_root, slug, &http_client, &config, &vocab, &events_tx) .await; // Assert: the transcript lands. Poll because the worker runs // asynchronously after the heal hands off the job. let transcript = case_dir.join("2026-04-15T20-30-00Z.transcript.txt"); let start = std::time::Instant::now(); while !transcript.exists() { if start.elapsed() > Duration::from_secs(10) { panic!( "transcript did not appear within 10s (whisper hits: {})", whisper.received_requests().await.unwrap().len() ); } tokio::time::sleep(Duration::from_millis(25)).await; } let body = std::fs::read_to_string(&transcript).unwrap(); assert_eq!(body, "recovered text"); // Teardown: close the worker channel so the background task exits. drop(tx_t); drop(pipeline); let _ = tokio::time::timeout(Duration::from_secs(5), worker).await; }