Files
doctate/server/tests/transient_failure_retries_test.rs
T
Brummel 3218fc62bd feat: server-authoritative case window via users.toml
The /api/oneliners time window is now read per-user from users.toml
(window_hours, default 72h). Clients no longer carry a window:
client.toml oneliner_window_hours, SyncConfig.window_hours, the
?hours=N query param, and the render_case_list cutoff filter are
gone. ETag suffix keeps the effective hours so an admin edit to
users.toml invalidates client caches on the next request.
OnelinersResponse.window_hours stays in the wire format, but now
exists solely to anchor client reconciliation.
2026-04-21 16:48:01 +02:00

174 lines
6.4 KiB
Rust

//! 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,
}];
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;
}