From 6e6a77fadbb1064a6b23dcb74401e38b9d5174b2 Mon Sep 17 00:00:00 2001 From: kurosawareiji7007-hub Date: Fri, 31 Jul 2026 13:01:04 -0700 Subject: [PATCH 1/3] fix(stt): bound whisper-server HTTP reads and fail closed on errors MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The STT server path still used unbounded read_line, had no connect/request deadline or status check, and fail-opened non-JSON bodies into fake transcripts — the same class of bug fixed for LLM/HA in #655/#656. --- CHANGELOG.md | 7 + crates/genie-core/src/voice/stt.rs | 206 ++++++++++++++++++++++------- 2 files changed, 165 insertions(+), 48 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index bce7c344..0c4ed222 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,13 @@ ## Unreleased +### Reliability + +- **STT (whisper-server)**: bound the HTTP response reader (shared + `read_response` limits + connect/request deadlines), require 2xx, and fail + closed on non-JSON bodies so a broken whisper-server cannot OOM the voice + loop or be spoken as a transcript (#936). + ### Quick-router tool-call accuracy - **Scene / routine**: strip a leading `please` before the exact-match set and diff --git a/crates/genie-core/src/voice/stt.rs b/crates/genie-core/src/voice/stt.rs index ca3d4648..15c9db21 100644 --- a/crates/genie-core/src/voice/stt.rs +++ b/crates/genie-core/src/voice/stt.rs @@ -1,8 +1,10 @@ use anyhow::Result; +use genie_common::http::{HttpReadError, HttpResponseLimits, read_response}; use std::sync::Mutex; use std::sync::atomic::{AtomicU64, Ordering}; use std::time::Duration; -use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; +use tokio::io::{AsyncWriteExt, BufReader}; +use tokio::net::TcpStream; use tokio::process::{Child, Command}; use tokio::sync::mpsc; @@ -11,6 +13,13 @@ use tokio::sync::mpsc; /// hung CUDA/NvMap child cannot wedge the single per-device voice loop forever. const WHISPER_CLI_TIMEOUT: Duration = Duration::from_secs(60); +/// Connect / full-request deadlines for whisper-server HTTP STT (#936). +const WHISPER_SERVER_CONNECT_TIMEOUT: Duration = Duration::from_secs(5); +const WHISPER_SERVER_REQUEST_TIMEOUT: Duration = Duration::from_secs(60); +/// Whisper JSON transcripts are small; a generous cap still blocks OOM from a +/// runaway/chunked body (same class as #656 for LLM/HA clients). +const WHISPER_SERVER_MAX_BODY_BYTES: usize = 256 * 1024; + /// Run `whisper-cli` under a deadline, reaping the child on timeout via /// `kill_on_drop`. async fn whisper_cli_output(mut cmd: Command) -> Result { @@ -342,61 +351,67 @@ impl SttEngine { let content_length = text_parts.len() + file_part.len() + wav_data.len() + body_end.len(); - let stream = tokio::net::TcpStream::connect(&addr).await?; + let stream = match tokio::time::timeout( + WHISPER_SERVER_CONNECT_TIMEOUT, + TcpStream::connect(&addr), + ) + .await + { + Ok(Ok(stream)) => stream, + Ok(Err(e)) => return Err(e.into()), + Err(_) => anyhow::bail!( + "whisper-server connect timed out after {}s", + WHISPER_SERVER_CONNECT_TIMEOUT.as_secs() + ), + }; let (reader, mut writer) = stream.into_split(); let request = format!( "POST /inference HTTP/1.1\r\nHost: {addr}\r\nContent-Type: multipart/form-data; boundary={boundary}\r\nContent-Length: {content_length}\r\nConnection: close\r\n\r\n" ); - writer.write_all(request.as_bytes()).await?; - writer.write_all(text_parts.as_bytes()).await?; - writer.write_all(file_part.as_bytes()).await?; - writer.write_all(&wav_data).await?; - writer.write_all(body_end.as_bytes()).await?; - - // Read response. - let mut buf_reader = BufReader::new(reader); - let mut response_body = String::new(); - let mut in_body = false; - - loop { - let mut line = String::new(); - let n = buf_reader.read_line(&mut line).await?; - if n == 0 { - break; + // Bound the write + response read the same way HA/LLM outbound clients + // do (#655/#656). The previous hand-rolled `read_line` loop had no + // line/body cap, no status check, and fail-opened non-JSON into a fake + // transcript (#936). + let language_hint = self.language_hint.clone(); + let response = match tokio::time::timeout(WHISPER_SERVER_REQUEST_TIMEOUT, async move { + writer.write_all(request.as_bytes()).await?; + writer.write_all(text_parts.as_bytes()).await?; + writer.write_all(file_part.as_bytes()).await?; + writer.write_all(&wav_data).await?; + writer.write_all(body_end.as_bytes()).await?; + + let mut buf_reader = BufReader::new(reader); + let limits = HttpResponseLimits::with_body_cap(WHISPER_SERVER_MAX_BODY_BYTES); + let response = read_response(&mut buf_reader, &limits) + .await + .map_err(|e| match e { + HttpReadError::BodyTooLarge => { + anyhow::anyhow!("whisper-server response exceeded body cap") + } + other => anyhow::anyhow!("whisper-server response read failed: {other}"), + })?; + if response.status == 0 { + anyhow::bail!("whisper-server returned an invalid HTTP status line"); } - if in_body { - response_body.push_str(&line); - } else if line.trim().is_empty() { - in_body = true; - } - } - - // Parse whisper server JSON response. - // Format: {"text": " Hello, turn on the lights."} - let parsed: serde_json::Value = serde_json::from_str(response_body.trim()) - .unwrap_or_else(|_| serde_json::json!({"text": response_body.trim()})); - - let text = Self::clean_hallucinations( - parsed - .get("text") - .and_then(|v| v.as_str()) - .unwrap_or("") - .trim(), - ); - let detected_language = super::language::detect_language_from_text(&text); - - Ok(Transcript { - text, - duration_ms: 0, - language: parsed - .get("language") - .and_then(|v| v.as_str()) - .map(String::from) - .or_else(|| self.language_hint.clone()) - .or(detected_language), + Ok::<_, anyhow::Error>(response) }) + .await + { + Ok(Ok(response)) => response, + Ok(Err(e)) => return Err(e), + Err(_) => anyhow::bail!( + "whisper-server request timed out after {}s", + WHISPER_SERVER_REQUEST_TIMEOUT.as_secs() + ), + }; + + transcript_from_whisper_http_response( + response.status, + &response.body, + language_hint.as_deref(), + ) } async fn transcribe_via_cli(&self, wav_path: &str) -> Result { @@ -559,6 +574,51 @@ impl SttEngine { } } +/// Decode a whisper-server HTTP response into a [`Transcript`]. +/// +/// Fail-closed on non-2xx and non-JSON bodies — the previous path stuffed the +/// raw body into `{"text": ...}`, so HTML error pages could be spoken (#936). +fn transcript_from_whisper_http_response( + status: u16, + body: &str, + language_hint: Option<&str>, +) -> Result { + if !(200..300).contains(&status) { + let trimmed = body.trim().replace('\n', " "); + anyhow::bail!( + "whisper-server HTTP {status}{}", + if trimmed.is_empty() { + String::new() + } else { + format!(": {trimmed}") + } + ); + } + + let parsed: serde_json::Value = serde_json::from_str(body.trim()) + .map_err(|e| anyhow::anyhow!("whisper-server returned non-JSON body: {e}"))?; + + let text = SttEngine::clean_hallucinations( + parsed + .get("text") + .and_then(|v| v.as_str()) + .unwrap_or("") + .trim(), + ); + let detected_language = super::language::detect_language_from_text(&text); + + Ok(Transcript { + text, + duration_ms: 0, + language: parsed + .get("language") + .and_then(|v| v.as_str()) + .map(String::from) + .or_else(|| language_hint.map(str::to_string)) + .or(detected_language), + }) +} + /// Drain stale samples from the ALSA capture queue before a real recording. /// Without this, between-cycle residue (kernel DMA carry-over from the prior /// arecord, plus a few hundred ms of acoustic echo from TTS playback bleeding @@ -1130,4 +1190,54 @@ mod tests { "what's the weather like" ); } + + #[test] + fn whisper_server_http_accepts_json_transcript() { + let t = transcript_from_whisper_http_response( + 200, + r#"{"text":" turn on the lights.","language":"en"}"#, + None, + ) + .unwrap(); + assert_eq!(t.text, "turn on the lights."); + assert_eq!(t.language.as_deref(), Some("en")); + } + + #[test] + fn whisper_server_http_rejects_non_2xx() { + let err = transcript_from_whisper_http_response(500, "error", None) + .unwrap_err() + .to_string(); + assert!(err.contains("HTTP 500"), "{err}"); + assert!(err.contains("html"), "{err}"); + } + + #[test] + fn whisper_server_http_rejects_non_json_body() { + // Previously fail-opened into a fake transcript of the HTML/plain body. + let err = transcript_from_whisper_http_response(200, "Internal Server Error", None) + .unwrap_err() + .to_string(); + assert!(err.contains("non-JSON"), "{err}"); + } + + #[tokio::test] + async fn whisper_server_read_rejects_oversized_header_line() { + // Prove the shared bounded reader is what STT now uses: an endless + // header line must fail closed instead of growing without bound. + let mut reader = BufReader::new(tokio::io::Cursor::new( + b"HTTP/1.1 200 OK\r\nX-Overflow: " + .iter() + .copied() + .chain(std::iter::repeat(b'a').take(9 * 1024)) + .chain(b"\r\n\r\n{}\r\n".iter().copied()) + .collect::>(), + )); + let limits = HttpResponseLimits::with_body_cap(WHISPER_SERVER_MAX_BODY_BYTES); + let err = read_response(&mut reader, &limits).await.unwrap_err(); + assert!( + matches!(err, HttpReadError::HeaderLineTooLong), + "got {err:?}" + ); + } } From 71d0ff96a6c8f034fdb588b7a162f92706d5f634 Mon Sep 17 00:00:00 2001 From: kurosawareiji7007-hub Date: Fri, 31 Jul 2026 14:33:41 -0700 Subject: [PATCH 2/3] fix(stt): rustfmt connect timeout and use std::io::Cursor in test --- crates/genie-core/src/voice/stt.rs | 27 +++++++++++++-------------- 1 file changed, 13 insertions(+), 14 deletions(-) diff --git a/crates/genie-core/src/voice/stt.rs b/crates/genie-core/src/voice/stt.rs index 15c9db21..e1afc9f8 100644 --- a/crates/genie-core/src/voice/stt.rs +++ b/crates/genie-core/src/voice/stt.rs @@ -351,19 +351,17 @@ impl SttEngine { let content_length = text_parts.len() + file_part.len() + wav_data.len() + body_end.len(); - let stream = match tokio::time::timeout( - WHISPER_SERVER_CONNECT_TIMEOUT, - TcpStream::connect(&addr), - ) - .await - { - Ok(Ok(stream)) => stream, - Ok(Err(e)) => return Err(e.into()), - Err(_) => anyhow::bail!( - "whisper-server connect timed out after {}s", - WHISPER_SERVER_CONNECT_TIMEOUT.as_secs() - ), - }; + let stream = + match tokio::time::timeout(WHISPER_SERVER_CONNECT_TIMEOUT, TcpStream::connect(&addr)) + .await + { + Ok(Ok(stream)) => stream, + Ok(Err(e)) => return Err(e.into()), + Err(_) => anyhow::bail!( + "whisper-server connect timed out after {}s", + WHISPER_SERVER_CONNECT_TIMEOUT.as_secs() + ), + }; let (reader, mut writer) = stream.into_split(); let request = format!( @@ -1067,6 +1065,7 @@ use std::sync::Arc; #[cfg(test)] mod tests { use super::*; + use std::io::Cursor; #[tokio::test] async fn write_wav_header() { @@ -1225,7 +1224,7 @@ mod tests { async fn whisper_server_read_rejects_oversized_header_line() { // Prove the shared bounded reader is what STT now uses: an endless // header line must fail closed instead of growing without bound. - let mut reader = BufReader::new(tokio::io::Cursor::new( + let mut reader = BufReader::new(Cursor::new( b"HTTP/1.1 200 OK\r\nX-Overflow: " .iter() .copied() From 0cf12f42a2f86a916ecc01f2595f2adc7ca27668 Mon Sep 17 00:00:00 2001 From: kurosawareiji7007-hub Date: Fri, 31 Jul 2026 15:03:08 -0700 Subject: [PATCH 3/3] fix(stt): use repeat_n in oversized-header test for clippy --- crates/genie-core/src/voice/stt.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/genie-core/src/voice/stt.rs b/crates/genie-core/src/voice/stt.rs index e1afc9f8..0bd9a28f 100644 --- a/crates/genie-core/src/voice/stt.rs +++ b/crates/genie-core/src/voice/stt.rs @@ -1228,7 +1228,7 @@ mod tests { b"HTTP/1.1 200 OK\r\nX-Overflow: " .iter() .copied() - .chain(std::iter::repeat(b'a').take(9 * 1024)) + .chain(std::iter::repeat_n(b'a', 9 * 1024)) .chain(b"\r\n\r\n{}\r\n".iter().copied()) .collect::>(), ));