Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions crates/tinyagents-harness/src/agent_loop/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
Expand Down
12 changes: 12 additions & 0 deletions crates/tinyagents-harness/src/agent_loop/model_call.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<State: Send + Sync, Ctx: Send + Sync> AgentHarness<State, Ctx> {
Expand Down Expand Up @@ -920,6 +921,7 @@ impl<State: Send + Sync, Ctx: Send + Sync> AgentHarness<State, Ctx> {
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
Expand Down Expand Up @@ -984,6 +986,10 @@ impl<State: Send + Sync, Ctx: Send + Sync> AgentHarness<State, Ctx> {
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
Expand Down Expand Up @@ -1053,6 +1059,12 @@ impl<State: Send + Sync, Ctx: Send + Sync> AgentHarness<State, Ctx> {
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()
Expand Down
28 changes: 28 additions & 0 deletions crates/tinyagents-harness/src/agent_loop/test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<String>();
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;
Expand Down
6 changes: 6 additions & 0 deletions crates/tinyagents-harness/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion crates/tinyagents-harness/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ pub use model_registry::{ModelRegistry, ModelSelection, ResolvedModelBinding};
pub use no_progress::{
ClassifiedFailure, ClassifiedFailureTracker, DEFAULT_IDENTICAL_HALT_THRESHOLD,
DEFAULT_REPEAT_CALL_THRESHOLD, DEFAULT_REPEAT_OUTPUT_THRESHOLD, NoProgress, NoProgressTracker,
SuccessfulRepeat, SuccessfulRepeatTracker, ToolAttempt,
StreamTextStallDetector, SuccessfulRepeat, SuccessfulRepeatTracker, ToolAttempt,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority high security confident

Wire the detector into streaming model calls

The new detector is only re-exported here; this change does not connect it to invoke_streaming or the streaming loop. Consequently, callers cannot receive the promised stall detection, and the earlier high-severity finding remains unresolved. Integrate the detector into each streaming model call before exposing the public export.

[RULE] dead-api ·

};
pub use observability::{
AgentCallLatency, AgentLatencyMetrics, AgentObservation, FanOutSink, HarnessEventJournal,
Expand Down
15 changes: 10 additions & 5 deletions crates/tinyagents-harness/src/no_progress/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -47,6 +48,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`).
Expand All @@ -57,6 +61,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. |

Expand Down
2 changes: 2 additions & 0 deletions crates/tinyagents-harness/src/no_progress/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,10 +61,12 @@
//! `as_str()` gives a stable telemetry label.

mod classified;
mod stream_text;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority critical security confident

Add the declared stream_text module

This module declaration requires no_progress/stream_text.rs or no_progress/stream_text/mod.rs, but neither exists in the reviewed tree. The crate therefore fails to compile until the module implementation is included or the declaration is removed.

[RULE] missing-module ·

mod successful_repeat;
mod types;

pub use classified::{ClassifiedFailure, ClassifiedFailureTracker};
pub use stream_text::StreamTextStallDetector;
Comment thread
senamakel marked this conversation as resolved.
pub use successful_repeat::{DEFAULT_REPEAT_CALL_THRESHOLD, DEFAULT_REPEAT_OUTPUT_THRESHOLD};
use types::LadderState;
pub use types::{
Expand Down
104 changes: 104 additions & 0 deletions crates/tinyagents-harness/src/no_progress/stream_text/mod.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
//! 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 {
Comment thread
senamakel marked this conversation as resolved.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

priority high security confident

Wire the detector into streaming model calls

This pull request only defines the detector; no streaming model-call path invokes observe, so the detector can never stop the open-ended narration described by the module. Connect it to the existing streaming delta loop, propagate a stalled result through the run's normal cancellation/error handling, and reset it at structured tool-call boundaries as intended.

[RULE] unwired-detector ·

sentence: String,
recent_starts: VecDeque<String>,
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 {
Comment thread
senamakel marked this conversation as resolved.
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;
self.stalled = false;
}

fn finish_sentence(&mut self) {
let words: Vec<String> = self
.sentence
.split_whitespace()
.map(|word| {
word.chars()
.filter(|ch| ch.is_alphanumeric())
.collect::<String>()
.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();
}
}

/// 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" | "ill now"
)
}

#[cfg(test)]
mod test;
118 changes: 118 additions & 0 deletions crates/tinyagents-harness/src/no_progress/stream_text/test.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
use super::*;

#[test]
fn catches_repeated_process_narration_across_chunks() {
let mut detector = StreamTextStallDetector::default();
Comment thread
senamakel marked this conversation as resolved.
let text = (0..12)
.map(|i| format!("Let me check the source number {i} before I answer the user. "))
.collect::<String>();
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() {
Comment thread
senamakel marked this conversation as resolved.
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!(
Comment thread
senamakel marked this conversation as resolved.
"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}. "
)));
}
}

#[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::<String>();
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::<String>();
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}. "
)));
}
}
Loading