From e23f4cb82e5b8562c97e8697e0a10a912ea5edfc Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 25 Sep 2026 17:47:37 +0530 Subject: [PATCH 1/4] Add streamed narration stall detector --- crates/tinyagents-harness/src/lib.rs | 2 +- .../src/no_progress/README.md | 4 + .../tinyagents-harness/src/no_progress/mod.rs | 2 + .../src/no_progress/stream_text/mod.rs | 91 +++++++++++++++++++ .../src/no_progress/stream_text/test.rs | 65 +++++++++++++ 5 files changed, 163 insertions(+), 1 deletion(-) create mode 100644 crates/tinyagents-harness/src/no_progress/stream_text/mod.rs create mode 100644 crates/tinyagents-harness/src/no_progress/stream_text/test.rs diff --git a/crates/tinyagents-harness/src/lib.rs b/crates/tinyagents-harness/src/lib.rs index 1474a22c6..ca329ce49 100644 --- a/crates/tinyagents-harness/src/lib.rs +++ b/crates/tinyagents-harness/src/lib.rs @@ -127,7 +127,7 @@ pub use model_registry::{ModelRegistry, ModelSelection, ResolvedModelBinding}; pub use no_progress::{ DEFAULT_IDENTICAL_HALT_THRESHOLD, DEFAULT_REPEAT_CALL_THRESHOLD, DEFAULT_REPEAT_OUTPUT_THRESHOLD, NoProgress, NoProgressTracker, SuccessfulRepeat, - SuccessfulRepeatTracker, ToolAttempt, + SuccessfulRepeatTracker, StreamTextStallDetector, ToolAttempt, }; pub use observability::{ AgentCallLatency, AgentLatencyMetrics, AgentObservation, FanOutSink, HarnessEventJournal, diff --git a/crates/tinyagents-harness/src/no_progress/README.md b/crates/tinyagents-harness/src/no_progress/README.md index ed90179b6..115c8cbd8 100644 --- a/crates/tinyagents-harness/src/no_progress/README.md +++ b/crates/tinyagents-harness/src/no_progress/README.md @@ -40,6 +40,9 @@ follow-up; see the "Driving this from an `after_tool` hook" section in - [`SuccessfulRepeatTracker`] / [`SuccessfulRepeat`] — the successful-repeat counterpart: `record_output`, `record_call_batch`, `record_call_outcome`, and `reset`. +- [`StreamTextStallDetector`] — consumes visible text fragments during one + model call and flags a long run of similarly opened sentences before the + provider stream finishes. - Threshold constants: [`DEFAULT_IDENTICAL_HALT_THRESHOLD`], [`DEFAULT_REPEAT_OUTPUT_THRESHOLD`], [`DEFAULT_REPEAT_CALL_THRESHOLD`] (all re-exported from `crate`). @@ -50,6 +53,7 @@ follow-up; see the "Driving this from an `after_tool` hook" section in | --- | --- | | `mod.rs` | The identical/any-failure escalation ladder (`NoProgressTracker::record`), argument fingerprinting, and the nudge/halt message builders. | | `successful_repeat.rs` | The successful-repeat streak tracker (`SuccessfulRepeatTracker`) and its private `Streak` helper. | +| `stream_text/` | Chunk-independent streamed-text stall detector and focused tests. | | `types.rs` | Public and crate-private type definitions shared by both trackers. | | `test.rs` | Unit tests for the escalation ladder. | diff --git a/crates/tinyagents-harness/src/no_progress/mod.rs b/crates/tinyagents-harness/src/no_progress/mod.rs index 2bf534374..09bb15075 100644 --- a/crates/tinyagents-harness/src/no_progress/mod.rs +++ b/crates/tinyagents-harness/src/no_progress/mod.rs @@ -60,9 +60,11 @@ //! and [`NoProgress::as_str`] exist so step 4 needs no enum match, and //! `as_str()` gives a stable telemetry label. +mod stream_text; mod successful_repeat; mod types; +pub use stream_text::StreamTextStallDetector; pub use successful_repeat::{DEFAULT_REPEAT_CALL_THRESHOLD, DEFAULT_REPEAT_OUTPUT_THRESHOLD}; use types::LadderState; pub use types::{ diff --git a/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs b/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs new file mode 100644 index 000000000..da4dd3ff4 --- /dev/null +++ b/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs @@ -0,0 +1,91 @@ +//! Detect a streamed response that keeps starting sentences the same way. +//! +//! The provider can stream an open-ended sequence of process narration without +//! ever completing a model call. Tool-call repeat guards run only after that +//! call, so they cannot stop this shape of stall. + +use std::collections::VecDeque; + +const MIN_CHARS: usize = 600; +const WINDOW: usize = 10; +const MATCHES_TO_STALL: usize = 8; + +/// Per-model-call detector for a long run of similarly opened sentences. +/// Feed visible text fragments in stream order, regardless of chunk boundaries. +#[derive(Default)] +pub struct StreamTextStallDetector { + sentence: String, + recent_starts: VecDeque, + total_chars: usize, +} + +impl StreamTextStallDetector { + /// Observe the next visible text fragment. Returns `true` once a strong + /// sentence-start recurrence is present; callers should stop that stream. + pub fn observe(&mut self, fragment: &str) -> bool { + for ch in fragment.chars() { + self.total_chars += 1; + if matches!(ch, '.' | '!' | '?' | '\n') { + self.finish_sentence(); + } else if self.sentence.len() < 1024 { + self.sentence.push(ch); + } + } + if self.total_chars < MIN_CHARS || self.recent_starts.len() < WINDOW { + return false; + } + self.recent_starts + .iter() + .filter(|start| is_process_start(start)) + .any(|candidate| { + self.recent_starts + .iter() + .filter(|start| *start == candidate) + .count() + >= MATCHES_TO_STALL + }) + } + + /// A structured tool call is progress; any later visible text starts a + /// fresh window instead of inheriting narration before that call. + pub fn reset(&mut self) { + self.sentence.clear(); + self.recent_starts.clear(); + self.total_chars = 0; + } + + fn finish_sentence(&mut self) { + let words: Vec<&str> = self.sentence.split_whitespace().collect(); + if words.len() >= 4 && self.sentence.chars().count() >= 24 { + let start = words[..2] + .iter() + .map(|word| { + word.chars() + .filter(|ch| ch.is_alphanumeric()) + .collect::() + .to_lowercase() + }) + .collect::>() + .join(" "); + if !start.trim().is_empty() { + self.recent_starts.push_back(start); + if self.recent_starts.len() > WINDOW { + self.recent_starts.pop_front(); + } + } + } + self.sentence.clear(); + } +} + +/// A repeated subject in an explanation ("Jev is ...") may be intentional. +/// The failure shape is repeated self-narration of actions never taken. +fn is_process_start(start: &str) -> bool { + matches!( + start, + "let me" | "i will" | "i need" | "i should" | "i can" | "ill now" + ) +} + +#[cfg(test)] +mod test; diff --git a/crates/tinyagents-harness/src/no_progress/stream_text/test.rs b/crates/tinyagents-harness/src/no_progress/stream_text/test.rs new file mode 100644 index 000000000..a0517452c --- /dev/null +++ b/crates/tinyagents-harness/src/no_progress/stream_text/test.rs @@ -0,0 +1,65 @@ +use super::*; + +#[test] +fn catches_repeated_process_narration_across_chunks() { + let mut detector = StreamTextStallDetector::default(); + let text = (0..12) + .map(|i| format!("Let me check the source number {i} before I answer the user. ")) + .collect::(); + let mut stalled = false; + for chunk in text.as_bytes().chunks(7) { + stalled |= detector.observe(std::str::from_utf8(chunk).unwrap()); + } + assert!(stalled); +} + +#[test] +fn varied_prose_and_short_repetitions_continue() { + let mut detector = StreamTextStallDetector::default(); + for sentence in [ + "Jev accepts typed questions about the supplied state.", + "TypeSafe describes Choice, Score and Noul answers.", + "A desktop host observes the screen before asking Jev.", + "The host executes a bounded action and verifies its result.", + ] { + assert!(!detector.observe(sentence)); + } + let mut detector = StreamTextStallDetector::default(); + for i in 0..7 { + assert!(!detector.observe(&format!( + "Let me check one more source number {i} before answering. " + ))); + } + let mut detector = StreamTextStallDetector::default(); + for i in 0..12 { + assert!(!detector.observe(&format!( + "Jev is a typed decision model used in example number {i}. " + ))); + } +} + +#[test] +fn tool_progress_resets_the_window() { + let mut detector = StreamTextStallDetector::default(); + for i in 0..7 { + detector.observe(&format!( + "Let me inspect the page number {i} before answering. " + )); + } + detector.reset(); + for i in 0..7 { + assert!(!detector.observe(&format!( + "Let me inspect the result number {i} before answering. " + ))); + } +} + +#[test] +fn useful_facts_between_planning_sentences_break_the_streak() { + let mut detector = StreamTextStallDetector::default(); + for i in 0..12 { + assert!(!detector.observe(&format!( + "Let me check source {i} before answering the user. Jev returns a typed Choice for state {i}. " + ))); + } +} From cbc857b69a23c7bbc2df01544e42b71f0c61e537 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 25 Sep 2026 17:49:23 +0530 Subject: [PATCH 2/4] Format stall detector export --- crates/tinyagents-harness/src/lib.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/tinyagents-harness/src/lib.rs b/crates/tinyagents-harness/src/lib.rs index ca329ce49..9e946c1ef 100644 --- a/crates/tinyagents-harness/src/lib.rs +++ b/crates/tinyagents-harness/src/lib.rs @@ -126,8 +126,8 @@ pub use ids::*; pub use model_registry::{ModelRegistry, ModelSelection, ResolvedModelBinding}; pub use no_progress::{ DEFAULT_IDENTICAL_HALT_THRESHOLD, DEFAULT_REPEAT_CALL_THRESHOLD, - DEFAULT_REPEAT_OUTPUT_THRESHOLD, NoProgress, NoProgressTracker, SuccessfulRepeat, - SuccessfulRepeatTracker, StreamTextStallDetector, ToolAttempt, + DEFAULT_REPEAT_OUTPUT_THRESHOLD, NoProgress, NoProgressTracker, StreamTextStallDetector, + SuccessfulRepeat, SuccessfulRepeatTracker, ToolAttempt, }; pub use observability::{ AgentCallLatency, AgentLatencyMetrics, AgentObservation, FanOutSink, HarnessEventJournal, From 981e9fb9e238da5f869bdc2b76b6a3ce848eee58 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 25 Sep 2026 18:02:19 +0530 Subject: [PATCH 3/4] Handle coalesced and short streamed narration --- .../src/no_progress/stream_text/mod.rs | 47 ++++++++++++------- .../src/no_progress/stream_text/test.rs | 43 +++++++++++++++++ 2 files changed, 73 insertions(+), 17 deletions(-) diff --git a/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs b/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs index da4dd3ff4..078e154cd 100644 --- a/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs +++ b/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs @@ -17,20 +17,33 @@ pub struct StreamTextStallDetector { sentence: String, recent_starts: VecDeque, total_chars: usize, + stalled: bool, } impl StreamTextStallDetector { /// Observe the next visible text fragment. Returns `true` once a strong /// sentence-start recurrence is present; callers should stop that stream. pub fn observe(&mut self, fragment: &str) -> bool { + if self.stalled { + return true; + } for ch in fragment.chars() { self.total_chars += 1; if matches!(ch, '.' | '!' | '?' | '\n') { self.finish_sentence(); + self.stalled |= self.has_stalled_window(); + if self.stalled { + return true; + } } else if self.sentence.len() < 1024 { self.sentence.push(ch); } } + self.stalled |= self.has_stalled_window(); + self.stalled + } + + fn has_stalled_window(&self) -> bool { if self.total_chars < MIN_CHARS || self.recent_starts.len() < WINDOW { return false; } @@ -52,26 +65,26 @@ impl StreamTextStallDetector { self.sentence.clear(); self.recent_starts.clear(); self.total_chars = 0; + self.stalled = false; } fn finish_sentence(&mut self) { - let words: Vec<&str> = self.sentence.split_whitespace().collect(); - if words.len() >= 4 && self.sentence.chars().count() >= 24 { - let start = words[..2] - .iter() - .map(|word| { - word.chars() - .filter(|ch| ch.is_alphanumeric()) - .collect::() - .to_lowercase() - }) - .collect::>() - .join(" "); - if !start.trim().is_empty() { - self.recent_starts.push_back(start); - if self.recent_starts.len() > WINDOW { - self.recent_starts.pop_front(); - } + let words: Vec = self + .sentence + .split_whitespace() + .map(|word| { + word.chars() + .filter(|ch| ch.is_alphanumeric()) + .collect::() + .to_lowercase() + }) + .filter(|word| !word.is_empty()) + .take(2) + .collect(); + if words.len() == 2 { + self.recent_starts.push_back(words.join(" ")); + if self.recent_starts.len() > WINDOW { + self.recent_starts.pop_front(); } } self.sentence.clear(); diff --git a/crates/tinyagents-harness/src/no_progress/stream_text/test.rs b/crates/tinyagents-harness/src/no_progress/stream_text/test.rs index a0517452c..b2032f130 100644 --- a/crates/tinyagents-harness/src/no_progress/stream_text/test.rs +++ b/crates/tinyagents-harness/src/no_progress/stream_text/test.rs @@ -63,3 +63,46 @@ fn useful_facts_between_planning_sentences_break_the_streak() { ))); } } + +#[test] +fn a_large_fragment_cannot_evict_an_earlier_stall() { + let repeated = (0..10) + .map(|i| { + format!( + "Let me check source {i} carefully before answering the user's request with the correct details. " + ) + }) + .collect::(); + let useful = "Jev returns typed choices. TypeSafe hosts its API. The caller verifies actions."; + let mut at_boundary = StreamTextStallDetector::default(); + let mut coalesced = StreamTextStallDetector::default(); + assert!(at_boundary.observe(&repeated)); + assert!(coalesced.observe(&format!("{repeated}{useful}"))); +} + +#[test] +fn short_narration_stalls_but_short_facts_break_the_window() { + let mut repeated = StreamTextStallDetector::default(); + assert!(!repeated.observe(&format!("{}.", "x".repeat(MIN_CHARS)))); + for _ in 0..9 { + assert!(!repeated.observe("Let me check. ")); + } + assert!(repeated.observe("Let me check.")); + + let mut interleaved = StreamTextStallDetector::default(); + for _ in 0..12 { + assert!( + !interleaved + .observe("Let me inspect this long source before answering. The answer is 42. ") + ); + } +} + +#[test] +fn markdown_list_markers_do_not_hide_narration() { + let mut detector = StreamTextStallDetector::default(); + let list = (0..12) + .map(|i| format!("- Let me inspect source {i} before answering the user.\n")) + .collect::(); + assert!(detector.observe(&list)); +} From ad9762f8596407394dd2baeb18cc728ee99102b4 Mon Sep 17 00:00:00 2001 From: Steven Enamakel Date: Fri, 25 Sep 2026 18:23:35 +0530 Subject: [PATCH 4/4] Stop stalled narration in the streaming agent loop --- .../src/agent_loop/README.md | 6 ++-- .../src/agent_loop/model_call.rs | 12 ++++++++ .../tinyagents-harness/src/agent_loop/test.rs | 28 +++++++++++++++++++ crates/tinyagents-harness/src/error.rs | 6 ++++ .../src/no_progress/README.md | 11 ++++---- .../src/no_progress/stream_text/mod.rs | 2 +- .../src/no_progress/stream_text/test.rs | 10 +++++++ 7 files changed, 67 insertions(+), 8 deletions(-) diff --git a/crates/tinyagents-harness/src/agent_loop/README.md b/crates/tinyagents-harness/src/agent_loop/README.md index 76203f3a3..2d1d18e86 100644 --- a/crates/tinyagents-harness/src/agent_loop/README.md +++ b/crates/tinyagents-harness/src/agent_loop/README.md @@ -148,8 +148,9 @@ it applies identically to the unary and streaming paths. `invoke_streaming_in_context[_with_status]` — the streaming counterparts of the `invoke*` family: each model call goes through `ChatModel::stream` instead of `ChatModel::invoke`, threading deltas through - every middleware's `on_model_delta` hook, but the loop still only returns - once the run is over. + every middleware's `on_model_delta` hook. Visible text also crosses the + streamed-text stall detector; repeated process narration drops the provider + stream and returns non-retryable `GenerationStalled`. - `AgentHarness::invoke_stream` / `invoke_stream_in_context` — a caller-facing event stream (`stream.rs`): yields every `AgentEvent` emitted during the run as `AgentStreamItem::Event`, then a single terminal @@ -162,6 +163,7 @@ it applies identically to the unary and streaming paths. `TinyAgentsError::LimitExceeded` (model/tool cap reached), `TinyAgentsError::Timeout` (wall-clock deadline elapsed), +`TinyAgentsError::GenerationStalled` (repetitive visible model stream), `TinyAgentsError::ModelNotFound` (no model resolvable), `TinyAgentsError::ToolNotFound` (model called an unregistered tool), or any error surfaced by a model, tool, middleware, or structured-output extraction. diff --git a/crates/tinyagents-harness/src/agent_loop/model_call.rs b/crates/tinyagents-harness/src/agent_loop/model_call.rs index 5318268d3..8b1758895 100644 --- a/crates/tinyagents-harness/src/agent_loop/model_call.rs +++ b/crates/tinyagents-harness/src/agent_loop/model_call.rs @@ -18,6 +18,7 @@ pub(super) const PER_CALL_BOUND_LABEL: &str = "per-model-call ceiling"; use super::*; use crate::cache::{CacheSkipReason, apply_prompt_cache_breakpoints, scoped_cache_key}; +use crate::no_progress::StreamTextStallDetector; use tinyinference_llm::cache::CachePolicy; impl AgentHarness { @@ -920,6 +921,7 @@ impl AgentHarness { let mut saw_streamed_content = false; let mut transformed_tools = StreamAccumulator::new(); let mut saw_tool_delta = false; + let mut stream_stall = StreamTextStallDetector::default(); // Some providers pad the very first streamed text chunk with // whitespace that is a wire-format artifact, not content (see @@ -984,6 +986,10 @@ impl AgentHarness { self.middleware .run_on_model_delta(ctx, state, &mut model_delta) .await?; + if stream_stall.observe(&model_delta.content) { + tracing::warn!("[stream] stopped repetitive model narration"); + return Err(TinyAgentsError::GenerationStalled); + } // Unconditional, not gated on the post-middleware content: // the pre-middleware tail here is always non-empty (the // surrounding `if` already checked it), matching the @@ -1053,6 +1059,12 @@ impl AgentHarness { self.middleware .run_on_model_delta(ctx, state, &mut model_delta) .await?; + if model_delta.tool_call.is_some() { + stream_stall.reset(); + } else if stream_stall.observe(&model_delta.content) { + tracing::warn!("[stream] stopped repetitive model narration"); + return Err(TinyAgentsError::GenerationStalled); + } saw_streamed_content |= !message_delta.text.is_empty() || !message_delta.reasoning.is_empty() || !model_delta.content.is_empty() diff --git a/crates/tinyagents-harness/src/agent_loop/test.rs b/crates/tinyagents-harness/src/agent_loop/test.rs index 75c4400aa..9d2114753 100644 --- a/crates/tinyagents-harness/src/agent_loop/test.rs +++ b/crates/tinyagents-harness/src/agent_loop/test.rs @@ -3337,6 +3337,34 @@ async fn invoke_streaming_fires_on_model_delta_per_delta_and_accumulates() { ); } +#[tokio::test] +async fn streamed_repetitive_narration_stops_the_live_model_call() { + use crate::testkit::StreamingMock; + + let repeated = (0..12) + .map(|i| { + format!("Let me inspect source {i} carefully before I answer the user's question. ") + }) + .collect::(); + let mut harness: AgentHarness<()> = AgentHarness::new(); + harness.register_model( + "stream", + Arc::new(StreamingMock::from_text_chunks([repeated])), + ); + + let error = harness + .invoke_streaming( + &(), + (), + RunConfig::new("stalled-stream"), + vec![Message::user("Find information about Jev")], + ) + .await + .expect_err("repetitive streamed narration must stop the run"); + assert!(matches!(&error, TinyAgentsError::GenerationStalled)); + assert!(!crate::retry::is_retryable(&error)); +} + #[tokio::test] async fn streaming_delta_transform_controls_final_run_and_cached_response() { use crate::cache::InMemoryResponseCache; diff --git a/crates/tinyagents-harness/src/error.rs b/crates/tinyagents-harness/src/error.rs index 1dcc6868d..3f2cdb3c0 100644 --- a/crates/tinyagents-harness/src/error.rs +++ b/crates/tinyagents-harness/src/error.rs @@ -81,6 +81,12 @@ pub enum TinyAgentsError { #[error("model error: {0}")] Model(String), + /// Visible model output became repetitive process narration during a + /// stream. The stream is dropped immediately; retrying the same request + /// would reproduce the stall, so this failure is terminal. + #[error("streamed response stalled on repeated narration")] + GenerationStalled, + /// A model provider call failed with the full structured detail preserved /// — HTTP status, provider error code, and whether retrying the same /// request may succeed — instead of flattened into a display string. diff --git a/crates/tinyagents-harness/src/no_progress/README.md b/crates/tinyagents-harness/src/no_progress/README.md index 115c8cbd8..eb853f2bc 100644 --- a/crates/tinyagents-harness/src/no_progress/README.md +++ b/crates/tinyagents-harness/src/no_progress/README.md @@ -19,11 +19,12 @@ A sibling tracker, [`SuccessfulRepeatTracker`], covers the complementary shape — a model that keeps *succeeding* at the same no-op call (or cycles through a short repeating sequence) without making progress. -Both trackers are deliberately free of harness types (no `RunContext`, no -`Message`) so they can be unit-tested in isolation. **Nothing in the crate -drives them yet** — wiring one into an `after_tool` middleware hook is a -follow-up; see the "Driving this from an `after_tool` hook" section in -`mod.rs` for the exact contract a driver must implement. +All trackers are free of harness types (no `RunContext`, no `Message`) so they +can be unit-tested in isolation. The agent loop drives +`StreamTextStallDetector` on visible streaming output and ends the model call +with non-retryable `GenerationStalled` when it fires. Tool-call trackers remain +host-wired through `after_tool`; see the "Driving this from an `after_tool` +hook" section in `mod.rs` for that contract. ## Public surface diff --git a/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs b/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs index 078e154cd..4dcf7950d 100644 --- a/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs +++ b/crates/tinyagents-harness/src/no_progress/stream_text/mod.rs @@ -96,7 +96,7 @@ impl StreamTextStallDetector { fn is_process_start(start: &str) -> bool { matches!( start, - "let me" | "i will" | "i need" | "i should" | "i can" | "ill now" + "let me" | "i will" | "i need" | "i should" | "ill now" ) } diff --git a/crates/tinyagents-harness/src/no_progress/stream_text/test.rs b/crates/tinyagents-harness/src/no_progress/stream_text/test.rs index b2032f130..91bd79852 100644 --- a/crates/tinyagents-harness/src/no_progress/stream_text/test.rs +++ b/crates/tinyagents-harness/src/no_progress/stream_text/test.rs @@ -106,3 +106,13 @@ fn markdown_list_markers_do_not_hide_narration() { .collect::(); assert!(detector.observe(&list)); } + +#[test] +fn ordinary_capability_sentences_do_not_count_as_process_narration() { + let mut detector = StreamTextStallDetector::default(); + for i in 0..12 { + assert!(!detector.observe(&format!( + "I can explain the typed decision result in useful detail number {i}. " + ))); + } +}