use pyo3::prelude::*;
use super::AgentId;
enum StepContent {
Output,
Error(String),
Empty,
}
const DEFAULT_MAX_CONSECUTIVE_MODEL_ERRORS: u32 = 3;
const DEFAULT_MAX_CONSECUTIVE_EMPTY_STEPS: u32 = 500;
#[derive(Debug, Clone, Copy)]
pub(crate) struct StreamLimits {
pub max_model_errors: u32,
pub max_empty_steps: u32,
pub channel_buffer: usize,
}
impl Default for StreamLimits {
fn default() -> Self {
Self {
max_model_errors: DEFAULT_MAX_CONSECUTIVE_MODEL_ERRORS,
max_empty_steps: DEFAULT_MAX_CONSECUTIVE_EMPTY_STEPS,
channel_buffer: crate::streaming::DEFAULT_CHANNEL_BUFFER,
}
}
}
impl StreamLimits {
pub fn from_config(config: &super::config::RuntimeConfig) -> Self {
Self {
max_model_errors: config
.max_consecutive_model_errors
.unwrap_or(DEFAULT_MAX_CONSECUTIVE_MODEL_ERRORS),
max_empty_steps: config
.max_consecutive_empty_steps
.unwrap_or(DEFAULT_MAX_CONSECUTIVE_EMPTY_STEPS),
channel_buffer: config
.streaming_channel_buffer
.unwrap_or(crate::streaming::DEFAULT_CHANNEL_BUFFER),
}
}
}
const RUNAWAY_THINKING_ERROR: &str = "aborted: model output contained only thinking with no text or tool calls \
after too many consecutive steps (runaway rumination)";
fn is_model_quality_error(message: &str) -> bool {
message.contains("model output")
}
#[derive(Default)]
struct StreamErrorState {
last_error: Option<String>,
output_after_error: bool,
consecutive_model_errors: u32,
consecutive_empty_steps: u32,
limits: StreamLimits,
}
impl StreamErrorState {
fn new(limits: StreamLimits) -> Self {
Self {
limits,
..Self::default()
}
}
fn observe(&mut self, content: &StepContent) -> bool {
match content {
StepContent::Error(msg) => {
self.consecutive_empty_steps = 0;
if is_model_quality_error(msg) {
self.consecutive_model_errors += 1;
} else {
self.consecutive_model_errors = 0;
}
self.last_error = Some(msg.clone());
self.output_after_error = false;
self.limits.max_model_errors > 0
&& self.consecutive_model_errors >= self.limits.max_model_errors
}
StepContent::Output => {
self.consecutive_model_errors = 0;
self.consecutive_empty_steps = 0;
if self.last_error.is_some() {
self.output_after_error = true;
}
false
}
StepContent::Empty => {
self.consecutive_empty_steps += 1;
if self.limits.max_empty_steps > 0
&& self.consecutive_empty_steps >= self.limits.max_empty_steps
{
self.last_error = Some(RUNAWAY_THINKING_ERROR.to_string());
self.output_after_error = false;
true
} else {
false
}
}
}
}
}
async fn forward_step_to_writer(
writer: &crate::streaming::ChatResponseWriter,
mut step: crate::types::Step,
agent_id: AgentId,
streamed_text: &mut String,
) -> StepContent {
let has_error_status = step.status == crate::types::StepStatus::Error;
let has_error_field = !step.error.is_empty();
if has_error_status || has_error_field {
let error_msg = format_error_message(&step);
let http_code = step.http_code;
let is_model_quality = is_model_quality_error(&error_msg);
tracing::warn!(
agent_id = ?agent_id,
http_code,
error = %error_msg,
"{}",
if is_model_quality {
"Model produced invalid output. Stream continues (backend will retry)"
} else {
"Error step received. Stream continues (backend controls iteration)"
}
);
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.step,
&writer.step_tx,
std::mem::take(&mut step),
"step",
)
.await;
return StepContent::Error(error_msg);
}
let step_idx = step.step_index;
let tool_names: Vec<String> = step.tool_calls.iter().map(|tc| tc.name.clone()).collect();
let usage_summary = step.usage_metadata.as_ref().map(|u| {
format!(
"{}p/{}o/{}t",
u.prompt_token_count.unwrap_or(0),
u.candidates_token_count.unwrap_or(0),
u.thoughts_token_count.unwrap_or(0),
)
});
let text_len = step.content.len() + step.content_delta.len();
let thinking_len = step.thinking.len() + step.thinking_delta.len();
let has_tool_calls = !step.tool_calls.is_empty();
let is_complete_response = step.is_complete_response == Some(true);
forward_text(writer, &mut step, streamed_text).await;
if is_complete_response {
streamed_text.clear();
}
forward_thoughts(writer, &mut step).await;
forward_tool_calls(writer, &mut step, agent_id).await;
apply_step_metadata(writer, &mut step);
crate::streaming::ChatResponseWriter::fan_out(&writer.subs.step, &writer.step_tx, step, "step")
.await;
if !tool_names.is_empty() {
tracing::info!(
agent_id = ?agent_id,
step = step_idx,
tools = ?tool_names,
usage = ?usage_summary,
"tool_call"
);
} else if text_len > 0 || thinking_len > 0 {
tracing::debug!(
agent_id = ?agent_id,
text_len,
thinking_len,
usage = ?usage_summary,
"model_output"
);
}
if text_len > 0 || has_tool_calls {
StepContent::Output
} else {
StepContent::Empty
}
}
fn format_error_message(step: &crate::types::Step) -> String {
if !step.error.is_empty() {
return step.error.clone();
}
let content = if step.content.is_empty() {
&step.content_delta
} else {
&step.content
};
format!("Step error (status={:?}): {content}", step.status)
}
async fn forward_text(
writer: &crate::streaming::ChatResponseWriter,
step: &mut crate::types::Step,
streamed_text: &mut String,
) {
let is_model = step.source == crate::types::StepSource::Model;
let (raw, is_delta) = if step.content_delta.is_empty() {
(std::mem::take(&mut step.content), false)
} else {
(std::mem::take(&mut step.content_delta), true)
};
if raw.is_empty() {
return;
}
let text = if is_model {
match dedup_model_text(raw, is_delta, streamed_text) {
Some(t) => t,
None => return,
}
} else {
raw
};
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.event,
&writer.event_tx,
crate::streaming::ResponseEvent::TextChunk(text.clone()),
"event",
)
.await;
if is_model {
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.chunk,
&writer.chunk_tx,
crate::streaming::StreamChunk::Text(text.clone()),
"chunk",
)
.await;
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.text,
&writer.text_tx,
text,
"text",
)
.await;
}
}
fn dedup_model_text(raw: String, is_delta: bool, streamed: &mut String) -> Option<String> {
if is_delta {
streamed.push_str(&raw);
return Some(raw);
}
if streamed.is_empty() {
streamed.push_str(&raw);
return Some(raw);
}
if raw == *streamed {
return None;
}
if let Some(suffix) = raw.strip_prefix(streamed.as_str()) {
let suffix = suffix.to_owned();
streamed.push_str(&suffix);
return Some(suffix);
}
streamed.clear();
streamed.push_str(&raw);
Some(raw)
}
async fn forward_thoughts(
writer: &crate::streaming::ChatResponseWriter,
step: &mut crate::types::Step,
) {
let thinking = if step.thinking_delta.is_empty() {
std::mem::take(&mut step.thinking)
} else {
std::mem::take(&mut step.thinking_delta)
};
if thinking.is_empty() {
return;
}
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.event,
&writer.event_tx,
crate::streaming::ResponseEvent::ThoughtChunk(thinking.clone()),
"event",
)
.await;
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.chunk,
&writer.chunk_tx,
crate::streaming::StreamChunk::Thought(thinking.clone()),
"chunk",
)
.await;
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.thought,
&writer.thought_tx,
thinking,
"thought",
)
.await;
}
async fn forward_tool_calls(
writer: &crate::streaming::ChatResponseWriter,
step: &mut crate::types::Step,
agent_id: AgentId,
) {
for tc in std::mem::take(&mut step.tool_calls) {
tracing::debug!(
agent_id = ?agent_id,
tool = %tc.name,
"Streaming tool call event"
);
let event = crate::streaming::ToolCallEvent {
name: tc.name,
args: tc.args,
id: tc.id,
canonical_path: tc.canonical_path,
};
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.event,
&writer.event_tx,
crate::streaming::ResponseEvent::ToolCall(event.clone()),
"event",
)
.await;
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.chunk,
&writer.chunk_tx,
crate::streaming::StreamChunk::ToolCall(event.clone()),
"chunk",
)
.await;
crate::streaming::ChatResponseWriter::fan_out(
&writer.subs.tool_call,
&writer.tool_call_tx,
event,
"tool_call",
)
.await;
}
}
fn apply_step_metadata(
writer: &crate::streaming::ChatResponseWriter,
step: &mut crate::types::Step,
) {
if let Some(usage) = step.usage_metadata.take() {
writer.set_usage(usage);
}
if let Some(out) = step.structured_output.take() {
writer.set_structured_output(out);
}
}
enum StepIterationResult {
Step(Box<crate::types::Step>),
Stop,
Error(String),
}
fn classify_py_step_error(err: &pyo3::PyErr, agent_id: AgentId) -> StepIterationResult {
let is_stop =
Python::attach(|py| err.is_instance_of::<pyo3::exceptions::PyStopAsyncIteration>(py));
if is_stop {
tracing::debug!(agent_id = ?agent_id, "Step stream ended (StopAsyncIteration)");
return StepIterationResult::Stop;
}
let err_msg = Python::attach(|py| crate::error::classify_py_error(py, err).to_string());
tracing::error!(agent_id = ?agent_id, error = %err_msg, "Python step iteration failed");
StepIterationResult::Error(err_msg)
}
async fn process_next_step_iteration(
aiter_py: &Py<PyAny>,
agent_id: AgentId,
) -> StepIterationResult {
let next_fut = Python::attach(|py| -> PyResult<_> {
let aiter_bound = aiter_py.bind(py);
let coro = aiter_bound.call_method0("__anext__")?;
pyo3_async_runtimes::tokio::into_future(coro)
});
let next_fut = match next_fut {
Ok(fut) => fut,
Err(e) => return classify_py_step_error(&e, agent_id),
};
let step_py = match next_fut.await {
Ok(obj) => obj,
Err(e) => return classify_py_step_error(&e, agent_id),
};
Python::attach(|py| {
let step_bound = step_py.bind(py);
if step_bound.is_none() {
return StepIterationResult::Stop;
}
match super::py_scripts::to_dict_py(step_bound)
.and_then(|d| d.extract::<crate::types::Step>())
{
Ok(step) => StepIterationResult::Step(Box::new(step)),
Err(e) => {
let err_msg = format!("Failed to extract Step from Python object: {e}");
tracing::error!(agent_id = ?agent_id, "{err_msg}");
StepIterationResult::Error(err_msg)
}
}
})
}
pub async fn stream_steps_to_writer(
writer: &crate::streaming::ChatResponseWriter,
agent_id: AgentId,
aiter_py: &Py<PyAny>,
limits: StreamLimits,
) {
tracing::debug!(agent_id = ?agent_id, ?limits, "Starting step streaming");
let mut state = StreamErrorState::new(limits);
let mut streamed_text = String::new();
loop {
match process_next_step_iteration(aiter_py, agent_id).await {
StepIterationResult::Step(step) => {
let content =
forward_step_to_writer(writer, *step, agent_id, &mut streamed_text).await;
if state.observe(&content) {
tracing::warn!(
agent_id = ?agent_id,
consecutive_model_errors = state.consecutive_model_errors,
consecutive_empty_steps = state.consecutive_empty_steps,
"Stopping stream (repeated invalid output or runaway \
thinking-only rumination) — handing off to orchestrator recovery"
);
break;
}
}
StepIterationResult::Stop => break,
StepIterationResult::Error(err_msg) => {
send_stream_error(writer, err_msg);
return;
}
}
}
if let Some(error_msg) = state.last_error {
if state.output_after_error {
tracing::info!(
agent_id = ?agent_id,
error = %error_msg,
"Stream recovered after error — not propagating"
);
} else {
tracing::warn!(
agent_id = ?agent_id,
error = %error_msg,
"Stream ended with unrecovered error — propagating"
);
send_stream_error(writer, error_msg);
}
}
}
fn send_stream_error(writer: &crate::streaming::ChatResponseWriter, message: String) {
if let Err(e) = writer
.error_tx
.try_send(crate::streaming::StreamError { message })
{
tracing::debug!("Error channel full or closed (first error wins): {e}");
}
}
#[cfg(test)]
mod tests {
use std::sync::atomic::Ordering;
use super::{AgentId, dedup_model_text, format_error_message, forward_step_to_writer};
use crate::types::{Step, StepSource, StepStatus};
fn step_with(status: StepStatus, error: &str, content: &str) -> Step {
Step {
status,
error: error.to_string(),
content: content.to_string(),
..Step::default()
}
}
fn model_delta(delta: &str) -> Step {
Step {
source: StepSource::Model,
content_delta: delta.to_string(),
..Step::default()
}
}
fn model_complete(content: &str) -> Step {
Step {
source: StepSource::Model,
content: content.to_string(),
is_complete_response: Some(true),
..Step::default()
}
}
async fn text_of(steps: Vec<Step>) -> String {
let (writer, handle) = crate::streaming::channel();
writer.subs.text.store(true, Ordering::Release);
let mut streamed = String::new();
for step in steps {
forward_step_to_writer(&writer, step, AgentId(1), &mut streamed).await;
}
drop(writer);
handle.text().await.expect("text drains cleanly").into()
}
#[tokio::test]
async fn consolidated_complete_response_is_not_double_emitted() {
let text = text_of(vec![
model_delta("Healthy mock response"),
model_complete("Healthy mock response"),
])
.await;
assert_eq!(text, "Healthy mock response");
}
#[tokio::test]
async fn incremental_deltas_concatenate_once() {
let text = text_of(vec![
model_delta("Heal"),
model_delta("thy "),
model_delta("mock "),
model_delta("response"),
model_complete("Healthy mock response"),
])
.await;
assert_eq!(text, "Healthy mock response");
}
#[tokio::test]
async fn non_streaming_single_content_step_emitted_once() {
let text = text_of(vec![model_complete("Only once")]).await;
assert_eq!(text, "Only once");
}
#[tokio::test]
async fn two_messages_each_emitted_once() {
let text = text_of(vec![
model_delta("one"),
model_complete("one"),
model_delta("two"),
model_complete("two"),
])
.await;
assert_eq!(text, "onetwo");
}
#[test]
fn dedup_model_text_skips_exact_consolidation() {
let mut s = String::new();
assert_eq!(
dedup_model_text("abc".to_owned(), true, &mut s),
Some("abc".to_owned())
);
assert_eq!(dedup_model_text("abc".to_owned(), false, &mut s), None);
}
#[test]
fn dedup_model_text_trims_grown_snapshot() {
let mut s = String::new();
assert_eq!(
dedup_model_text("ab".to_owned(), true, &mut s),
Some("ab".to_owned())
);
assert_eq!(
dedup_model_text("abcd".to_owned(), false, &mut s),
Some("cd".to_owned())
);
}
#[test]
fn dedup_model_text_non_streaming_emits_content() {
let mut s = String::new();
assert_eq!(
dedup_model_text("full".to_owned(), false, &mut s),
Some("full".to_owned())
);
}
#[test]
fn error_status_is_detected() {
let step = step_with(StepStatus::Error, "", "some content");
assert_eq!(step.status, StepStatus::Error);
assert!(step.error.is_empty());
let has_error_status = step.status == StepStatus::Error;
let has_error_field = !step.error.is_empty();
assert!(has_error_status || has_error_field);
}
#[test]
fn error_field_is_detected() {
let step = step_with(StepStatus::Done, "quota exceeded", "");
let has_error_status = step.status == StepStatus::Error;
let has_error_field = !step.error.is_empty();
assert!(has_error_status || has_error_field);
}
#[test]
fn both_error_signals_detected() {
let step = step_with(StepStatus::Error, "model not found", "error text");
let has_error_status = step.status == StepStatus::Error;
let has_error_field = !step.error.is_empty();
assert!(has_error_status && has_error_field);
}
#[test]
fn normal_step_not_treated_as_error() {
let step = step_with(StepStatus::Done, "", "normal content");
let has_error_status = step.status == StepStatus::Error;
let has_error_field = !step.error.is_empty();
assert!(!has_error_status && !has_error_field);
}
#[test]
fn empty_content_with_done_status_is_not_error() {
let step = step_with(StepStatus::Done, "", "");
let has_error_status = step.status == StepStatus::Error;
let has_error_field = !step.error.is_empty();
assert!(!has_error_status && !has_error_field);
}
#[test]
fn format_uses_error_field_when_present() {
let step = step_with(StepStatus::Error, "quota exceeded", "some content");
assert_eq!(format_error_message(&step), "quota exceeded");
}
#[test]
fn format_falls_back_to_content_when_no_error_field() {
let step = step_with(StepStatus::Error, "", "agent terminated");
let msg = format_error_message(&step);
assert!(msg.contains("agent terminated"), "got: {msg}");
assert!(msg.contains("Error"), "got: {msg}");
}
#[test]
fn format_uses_content_delta_when_content_empty() {
let step = Step {
status: StepStatus::Error,
content_delta: "delta error text".to_string(),
..Step::default()
};
let msg = format_error_message(&step);
assert!(msg.contains("delta error text"), "got: {msg}");
}
fn simulate_stream(events: &[super::StepContent]) -> (Option<String>, bool) {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
for event in events {
state.observe(event);
}
(state.last_error, state.output_after_error)
}
#[test]
fn error_only_propagates() {
let (last_error, output_after) =
simulate_stream(&[super::StepContent::Error("503 unavailable".into())]);
assert!(last_error.is_some());
assert!(!output_after, "No output after error → should propagate");
}
#[test]
fn error_then_output_is_recovered() {
let (last_error, output_after) = simulate_stream(&[
super::StepContent::Error("model output empty".into()),
super::StepContent::Output,
]);
assert!(last_error.is_some());
assert!(
output_after,
"Output after error → recovered, don't propagate"
);
}
#[test]
fn error_then_output_then_error_propagates() {
let (last_error, output_after) = simulate_stream(&[
super::StepContent::Error("first error".into()),
super::StepContent::Output,
super::StepContent::Error("second error".into()),
]);
assert_eq!(last_error.as_deref(), Some("second error"));
assert!(!output_after, "Last error had no output after → propagate");
}
#[test]
fn clean_stream_no_error() {
let (last_error, output_after) = simulate_stream(&[
super::StepContent::Output,
super::StepContent::Empty,
super::StepContent::Output,
]);
assert!(last_error.is_none());
assert!(!output_after);
}
#[test]
fn output_before_error_does_not_count_as_recovery() {
let (last_error, output_after) = simulate_stream(&[
super::StepContent::Output, super::StepContent::Error("late error".into()),
]);
assert!(last_error.is_some());
assert!(
!output_after,
"Output before (not after) error → should propagate"
);
}
#[test]
fn empty_steps_do_not_affect_recovery() {
let (last_error, output_after) = simulate_stream(&[
super::StepContent::Error("error".into()),
super::StepContent::Empty,
super::StepContent::Empty,
]);
assert!(last_error.is_some());
assert!(!output_after, "Empty steps don't count as recovery");
}
#[test]
fn multiple_errors_then_output_is_recovered() {
let (last_error, output_after) = simulate_stream(&[
super::StepContent::Error("first".into()),
super::StepContent::Error("second".into()),
super::StepContent::Output,
]);
assert_eq!(last_error.as_deref(), Some("second"));
assert!(output_after, "Output after last error → recovered");
}
fn model_error() -> super::StepContent {
super::StepContent::Error(
"model output must contain either output text or tool calls".into(),
)
}
#[test]
fn three_consecutive_model_errors_stop_the_stream() {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
assert!(
!state.observe(&model_error()),
"1st model error keeps streaming"
);
assert!(
!state.observe(&model_error()),
"2nd model error keeps streaming"
);
assert!(
state.observe(&model_error()),
"3rd consecutive model error must stop the stream"
);
assert!(state.last_error.is_some());
assert!(
!state.output_after_error,
"no output followed → the error must propagate"
);
}
#[test]
fn output_resets_the_model_error_streak() {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
assert!(!state.observe(&model_error()));
assert!(!state.observe(&model_error()));
assert!(!state.observe(&super::StepContent::Output));
assert!(!state.observe(&model_error()));
assert!(!state.observe(&model_error()));
assert_eq!(state.consecutive_model_errors, 2);
}
#[test]
fn transient_errors_do_not_count_toward_the_model_limit() {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
for _ in 0..5 {
assert!(!state.observe(&super::StepContent::Error("503 unavailable".into())));
}
assert_eq!(state.consecutive_model_errors, 0);
assert!(state.last_error.is_some());
}
#[test]
fn a_transient_error_resets_the_model_error_streak() {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
assert!(!state.observe(&model_error()));
assert!(!state.observe(&model_error()));
assert!(!state.observe(&super::StepContent::Error("503 unavailable".into())));
assert_eq!(state.consecutive_model_errors, 0);
}
#[test]
fn empty_steps_do_not_reset_the_model_error_streak() {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
assert!(!state.observe(&model_error()));
assert!(!state.observe(&super::StepContent::Empty));
assert!(!state.observe(&model_error()));
assert!(state.observe(&model_error()));
}
#[test]
fn runaway_thinking_only_stream_is_aborted() {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
let limit = super::DEFAULT_MAX_CONSECUTIVE_EMPTY_STEPS;
for i in 1..limit {
assert!(
!state.observe(&super::StepContent::Empty),
"empty step {i} should not yet abort"
);
}
assert!(
state.observe(&super::StepContent::Empty),
"reaching the empty-step limit must abort the stream"
);
let err = state.last_error.expect("synthetic error recorded");
assert!(
super::is_model_quality_error(&err),
"synthetic runaway error must be model-quality so it routes to recovery"
);
assert!(!state.output_after_error, "no output → error propagates");
}
#[test]
fn output_resets_the_empty_step_streak() {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
for _ in 0..(super::DEFAULT_MAX_CONSECUTIVE_EMPTY_STEPS - 1) {
assert!(!state.observe(&super::StepContent::Empty));
}
assert!(!state.observe(&super::StepContent::Output));
assert_eq!(state.consecutive_empty_steps, 0);
assert!(!state.observe(&super::StepContent::Empty));
}
#[test]
fn interleaved_output_prevents_runaway_abort() {
let mut state = super::StreamErrorState::new(super::StreamLimits::default());
for _ in 0..10 {
for _ in 0..(super::DEFAULT_MAX_CONSECUTIVE_EMPTY_STEPS - 1) {
assert!(!state.observe(&super::StepContent::Empty));
}
assert!(!state.observe(&super::StepContent::Output));
}
assert!(state.last_error.is_none(), "healthy turn records no error");
}
#[test]
fn zero_model_error_limit_never_aborts() {
let limits = super::StreamLimits {
max_model_errors: 0,
max_empty_steps: super::DEFAULT_MAX_CONSECUTIVE_EMPTY_STEPS,
channel_buffer: crate::streaming::DEFAULT_CHANNEL_BUFFER,
};
let mut state = super::StreamErrorState::new(limits);
for _ in 0..100 {
assert!(
!state.observe(&model_error()),
"zero limit must never abort on model errors"
);
}
assert!(state.last_error.is_some());
}
#[test]
fn zero_empty_step_limit_never_aborts() {
let limits = super::StreamLimits {
max_model_errors: super::DEFAULT_MAX_CONSECUTIVE_MODEL_ERRORS,
max_empty_steps: 0,
channel_buffer: crate::streaming::DEFAULT_CHANNEL_BUFFER,
};
let mut state = super::StreamErrorState::new(limits);
for _ in 0..1000 {
assert!(
!state.observe(&super::StepContent::Empty),
"zero limit must never abort on empty steps"
);
}
assert!(state.last_error.is_none(), "no synthetic error recorded");
}
#[test]
fn stream_limits_from_config_uses_overrides() {
let config = super::super::config::RuntimeConfig {
max_consecutive_model_errors: Some(10),
max_consecutive_empty_steps: Some(42),
..Default::default()
};
let limits = super::StreamLimits::from_config(&config);
assert_eq!(limits.max_model_errors, 10);
assert_eq!(limits.max_empty_steps, 42);
}
#[test]
fn stream_limits_from_config_uses_defaults_for_none() {
let config = super::super::config::RuntimeConfig::default();
let limits = super::StreamLimits::from_config(&config);
assert_eq!(
limits.max_model_errors,
super::DEFAULT_MAX_CONSECUTIVE_MODEL_ERRORS
);
assert_eq!(
limits.max_empty_steps,
super::DEFAULT_MAX_CONSECUTIVE_EMPTY_STEPS
);
}
}