use std::future::Future;
use std::path::PathBuf;
use std::sync::Arc;
use crate::event_log::EventLog;
use crate::value::{VmError, VmValue};
use super::api::{
vm_call_llm_full_single_route, vm_call_llm_full_streaming_offthread_single_route,
vm_call_llm_full_streaming_single_route, DeltaSender,
};
use super::trace::{trace_llm_call, LlmTraceEntry};
use super::{api, cost, first_token, rate_limit, resolved_dispatch, routing, trace};
use super::agent_tools::next_call_id;
mod raw_tool_receipts;
mod served_context_receipts;
mod transcript_ambient;
use transcript_ambient::{current_transcript_dir, system_prompt_changed, tool_schemas_changed};
pub(crate) use transcript_ambient::{
pop_llm_transcript_dir, push_llm_transcript_dir, swap_llm_transcript_ambient,
LlmTranscriptAmbient,
};
tokio::task_local! {
static RAW_PROVIDER_CAPTURE_CONTEXT: RawProviderCaptureContext;
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct RawProviderCaptureContext {
pub(crate) call_id: String,
pub(crate) iteration: usize,
transcript_dir: Option<String>,
}
impl RawProviderCaptureContext {
fn new(call_id: String, iteration: usize) -> Self {
Self {
call_id,
iteration,
transcript_dir: current_transcript_dir(),
}
}
}
pub(crate) async fn with_raw_provider_capture_context<F>(
context: RawProviderCaptureContext,
future: F,
) -> F::Output
where
F: Future,
{
RAW_PROVIDER_CAPTURE_CONTEXT.scope(context, future).await
}
pub(crate) fn current_raw_provider_capture_context() -> Option<RawProviderCaptureContext> {
RAW_PROVIDER_CAPTURE_CONTEXT.try_with(Clone::clone).ok()
}
fn hash_str(value: &str) -> u64 {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
value.hash(&mut hasher);
hasher.finish()
}
fn hash_json(value: &serde_json::Value) -> u64 {
let encoded = serde_json::to_string(value).unwrap_or_default();
hash_str(&encoded)
}
fn env_flag_enabled(name: &str) -> bool {
std::env::var(name)
.ok()
.map(|value| {
let normalized = value.trim().to_ascii_lowercase();
matches!(normalized.as_str(), "1" | "true" | "yes" | "on" | "full")
})
.unwrap_or(false)
}
fn verbose_llm_transcript_enabled() -> bool {
env_flag_enabled("HARN_LLM_TRANSCRIPT_VERBOSE")
}
fn raw_llm_transcript_enabled() -> bool {
env_flag_enabled("HARN_LLM_TRANSCRIPT_RAW")
}
pub(crate) fn raw_provider_capture_enabled(context: Option<&RawProviderCaptureContext>) -> bool {
context
.and_then(|context| context.transcript_dir.as_ref())
.is_some()
&& raw_llm_transcript_enabled()
}
pub(crate) fn persist_raw_provider_request(
context: Option<&RawProviderCaptureContext>,
provider: &str,
model: &str,
wire_dialect: &str,
attempt: Option<usize>,
body: &serde_json::Value,
) -> Option<String> {
let context = context?;
if !raw_provider_capture_enabled(Some(context)) {
return None;
}
let envelope = serde_json::json!({
"schema_version": "harn.llm.raw_provider_request.v1",
"kind": "request",
"captured_at": chrono_now(),
"call_id": context.call_id,
"iteration": context.iteration,
"attempt": attempt,
"provider": provider,
"model": model,
"wire_dialect": wire_dialect,
"body": body,
});
write_raw_provider_sidecar(context, "request", provider, model, attempt, envelope)
}
pub(crate) fn persist_raw_provider_response(
context: Option<&RawProviderCaptureContext>,
provider: &str,
model: &str,
transport: &str,
attempt: Option<usize>,
status: u16,
content_type: Option<&str>,
body_text: &str,
) -> Option<String> {
let context = context?;
if !raw_provider_capture_enabled(Some(context)) {
return None;
}
let parsed_json = serde_json::from_str::<serde_json::Value>(body_text).ok();
let envelope = serde_json::json!({
"schema_version": "harn.llm.raw_provider_response.v1",
"kind": "response",
"captured_at": chrono_now(),
"call_id": context.call_id,
"iteration": context.iteration,
"attempt": attempt,
"provider": provider,
"model": model,
"transport": transport,
"status": status,
"content_type": content_type,
"body_text": body_text,
"body_json": parsed_json,
});
write_raw_provider_sidecar(
context,
&format!("response-{transport}"),
provider,
model,
attempt,
envelope,
)
}
fn write_raw_provider_sidecar(
context: &RawProviderCaptureContext,
suffix: &str,
provider: &str,
model: &str,
attempt: Option<usize>,
mut envelope: serde_json::Value,
) -> Option<String> {
crate::redact::current_policy().redact_json_in_place(&mut envelope);
let dir = context.transcript_dir.as_deref()?;
let raw_dir = PathBuf::from(&dir).join("raw-provider");
std::fs::create_dir_all(&raw_dir).ok()?;
let call_id = raw_provider_file_id(&context.call_id);
let attempt_part = attempt
.map(|attempt| format!("-attempt-{attempt}"))
.unwrap_or_default();
let filename = format!("{call_id}{attempt_part}-{suffix}.json");
let relative_path = format!("raw-provider/{filename}");
let path = raw_dir.join(filename);
let encoded = serde_json::to_vec_pretty(&envelope).ok()?;
static RAW_PROVIDER_WRITE_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
{
let _guard = RAW_PROVIDER_WRITE_LOCK
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let mut file = std::fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(path)
.ok()?;
use std::io::Write;
file.write_all(&encoded).ok()?;
file.write_all(b"\n").ok()?;
}
append_llm_transcript_entry_to_dir(
&serde_json::json!({
"type": "provider_raw_capture",
"timestamp": chrono_now(),
"span_id": crate::tracing::current_span_id(),
"call_id": context.call_id,
"iteration": context.iteration,
"attempt": attempt,
"provider": provider,
"model": model,
"capture": suffix,
"path": relative_path,
}),
Some(dir),
);
Some(relative_path)
}
fn raw_provider_file_id(value: &str) -> String {
let sanitized: String = value
.chars()
.map(|ch| {
if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
ch
} else {
'_'
}
})
.collect();
if sanitized.is_empty() {
"call".to_string()
} else {
sanitized
}
}
pub(super) fn is_retryable_llm_error(err: &VmError) -> bool {
use crate::value::{classify_error_message, ErrorCategory};
let msg = match err {
VmError::CategorizedError { category, message } => {
let llm_info = crate::llm::api::classify_llm_error(category.clone(), message);
return if llm_info.reason == crate::llm::api::LlmErrorReason::Unknown {
category.is_transient()
} else {
llm_info.kind == crate::llm::api::LlmErrorKind::Transient
};
}
VmError::Thrown(crate::value::VmValue::Dict(d)) => {
if let Some(kind) = d.get("kind").map(|v| v.display()) {
return kind == "transient";
}
if let Some(category) = d.get("category").map(|v| v.display()) {
return ErrorCategory::parse(&category).is_transient();
}
return false;
}
VmError::Thrown(crate::value::VmValue::String(s)) => s.as_ref(),
VmError::Runtime(s) => s.as_str(),
_ => return false,
};
let category = classify_error_message(msg);
let llm_info = crate::llm::api::classify_llm_error(category, msg);
if llm_info.kind == crate::llm::api::LlmErrorKind::Transient {
return true;
}
if llm_info.reason != crate::llm::api::LlmErrorReason::Unknown {
return false;
}
let derived = classify_error_message(msg);
if derived != ErrorCategory::Generic {
return derived.is_transient();
}
let lower = msg.to_lowercase();
lower.contains("too many requests")
|| lower.contains("rate limit")
|| lower.contains("overloaded")
|| lower.contains("service unavailable")
|| lower.contains("bad gateway")
|| lower.contains("gateway timeout")
|| lower.contains("timed out")
|| lower.contains("timeout")
|| lower.contains("delivered no content")
|| lower.contains("eof")
}
pub(super) fn is_network_failure_llm_error(err: &VmError) -> bool {
let (category, message) = match err {
VmError::CategorizedError { category, message } => (category.clone(), message.clone()),
VmError::Thrown(crate::value::VmValue::String(s)) => {
(crate::value::classify_error_message(s), s.to_string())
}
VmError::Runtime(s) => (crate::value::classify_error_message(s), s.clone()),
_ => return false,
};
let reason = crate::llm::api::classify_llm_error(category, &message).reason;
matches!(
reason,
crate::llm::api::LlmErrorReason::NetworkError | crate::llm::api::LlmErrorReason::Timeout
)
}
pub(super) fn is_overloaded_llm_error(err: &VmError) -> bool {
crate::value::error_to_category(err) == crate::value::ErrorCategory::Overloaded
}
mod detector;
mod provider_errors;
mod transcript_observability;
use detector::*;
use provider_errors::*;
use transcript_observability::*;
pub(crate) use transcript_observability::{append_llm_observability_entry, record_template_render};
pub(super) fn extract_retry_after_ms(err: &VmError) -> Option<u64> {
let msg = match err {
VmError::Thrown(crate::value::VmValue::String(s)) => s.as_ref(),
VmError::Thrown(crate::value::VmValue::Dict(d)) => {
return d.get("retry_after_ms").and_then(|v| match v {
crate::value::VmValue::Int(ms) if *ms >= 0 => Some(*ms as u64),
_ => None,
});
}
VmError::CategorizedError { message, .. } => message.as_str(),
VmError::Runtime(s) => s.as_str(),
_ => return None,
};
parse_retry_after(msg)
}
pub(crate) fn parse_retry_after(msg: &str) -> Option<u64> {
const MAX_MS: u64 = 60_000;
let lower = msg.to_lowercase();
let pos = lower.find("retry-after:")?;
let after = &msg[pos + "retry-after:".len()..];
let end = after.find(['\r', '\n']).unwrap_or(after.len());
let value = after[..end].trim();
if value.is_empty() {
return None;
}
let numeric_prefix = value
.chars()
.take_while(|ch| ch.is_ascii_digit() || *ch == '.')
.collect::<String>();
if !numeric_prefix.is_empty() {
if let Ok(secs) = numeric_prefix.parse::<f64>() {
if !secs.is_finite() || secs < 0.0 {
return Some(0);
}
let ms = (secs * 1000.0) as u64;
return Some(ms.min(MAX_MS));
}
}
if let Ok(target) = httpdate::parse_http_date(value) {
let now = std::time::SystemTime::now();
let delta = target
.duration_since(now)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
return Some(delta.min(MAX_MS));
}
None
}
fn empty_completion_retry_budget(provider: &str) -> usize {
if crate::llm::providers::MockProvider::should_intercept(provider) {
0
} else {
EMPTY_COMPLETION_BUILTIN_RETRIES
}
}
fn llm_retry_backoff_ms(error: &VmError, attempt: usize, provider: &str) -> u64 {
if crate::llm::providers::MockProvider::should_intercept(provider) {
return 0;
}
match extract_retry_after_ms(error) {
Some(retry_after_ms) => retry_after_ms.saturating_add(retry_after_jitter_ms()),
None => base_retry_backoff_ms(attempt),
}
}
fn base_retry_backoff_ms(attempt: usize) -> u64 {
let ceil = DEFAULT_LLM_CALL_BACKOFF_MS.saturating_mul(1 << attempt.min(4));
equal_jitter_ms(ceil, &mut rand::rng())
}
fn equal_jitter_ms<R: rand::RngExt>(ceil: u64, rng: &mut R) -> u64 {
let half = ceil / 2;
if half == 0 {
return ceil;
}
half + rand_range_inclusive(half, rng)
}
fn retry_after_jitter_ms() -> u64 {
rand_range_inclusive(DEFAULT_LLM_CALL_BACKOFF_MS, &mut rand::rng())
}
fn rand_range_inclusive<R: rand::RngExt>(max: u64, rng: &mut R) -> u64 {
rng.random_range(0..max.saturating_add(1))
}
fn degrade_options_to_text_channel(
opts: &super::api::LlmCallOptions,
) -> super::api::LlmCallOptions {
let mut degraded = opts.clone();
degraded.native_tools = None;
degraded.output_format = super::api::OutputFormat::Text;
degraded.response_format = None;
degraded.json_schema = None;
degraded
}
fn degrade_options_to_non_streaming_transport(
opts: &super::api::LlmCallOptions,
) -> super::api::LlmCallOptions {
let mut degraded = opts.clone();
degraded.stream = false;
degraded
}
pub(crate) async fn observed_llm_call(
opts: &super::api::LlmCallOptions,
tool_format: Option<&str>,
bridge: Option<&Arc<crate::bridge::HostBridge>>,
iteration: Option<usize>,
user_visible: bool,
offthread: bool,
streaming_detector: Option<StreamingDetectorContext>,
delta_sink: Option<DeltaSender>,
) -> Result<super::api::LlmResult, VmError> {
let _in_flight_guard = super::call::InFlightLlmCallGuard::enter(opts);
let mut effective_tool_format = tool_format
.map(str::to_string)
.or_else(|| {
std::env::var("HARN_AGENT_TOOL_FORMAT")
.ok()
.filter(|value| !value.trim().is_empty())
})
.unwrap_or_else(|| crate::llm_config::default_tool_format(&opts.model, &opts.provider));
let mut working: std::borrow::Cow<'_, super::api::LlmCallOptions> =
std::borrow::Cow::Borrowed(opts);
let mut degraded_to_text = false;
let mut degraded_stream_transport = false;
let mut attempt = 0usize;
let mut empty_completion_retries = 0usize;
loop {
let opts: &super::api::LlmCallOptions = working.as_ref();
super::rate_limit::check_network_breaker_for_llm_call(opts)?;
let rate_limit_permit = super::rate_limit::acquire_permit_for_llm_call(opts).await?;
let governor_org_key = crate::llm::rate_governor::org_key_id(&opts.api_key);
let governor_est_tokens = governor_estimated_tokens(opts);
let governor_reserved =
await_governor_admission(&opts.provider, &governor_org_key, governor_est_tokens).await;
let provider_was_throttled_during_call =
crate::llm::rate_governor::provider_already_throttled(
&opts.provider,
&governor_org_key,
);
let call_id = next_call_id();
let prompt_chars: usize = opts
.messages
.iter()
.filter_map(|m| m.get("content").and_then(|c| c.as_str()))
.map(|s| s.len())
.sum();
let mut span_meta = vec![
("call_id", serde_json::json!(call_id.clone())),
("model", serde_json::json!(opts.model.clone())),
("provider", serde_json::json!(opts.provider.clone())),
("prompt_chars", serde_json::json!(prompt_chars)),
(
"route_policy",
serde_json::json!(opts.route_policy.as_label()),
),
(
"fallback_chain",
serde_json::json!(opts.fallback_chain.clone()),
),
];
if let Some(decision) = opts.routing_decision.as_ref() {
span_meta.push(("routing_decision", serde_json::json!(decision)));
}
if let Some(iter) = iteration {
span_meta.push(("iteration", serde_json::json!(iter)));
span_meta.push(("llm_attempt", serde_json::json!(attempt)));
}
annotate_current_span(&span_meta);
let mut call_start_meta =
serde_json::json!({"model": opts.model, "prompt_chars": prompt_chars});
call_start_meta["stream_publicly"] =
serde_json::json!(opts.response_format.as_deref() != Some("json"));
call_start_meta["user_visible"] = serde_json::json!(user_visible);
if let Some(iter) = iteration {
call_start_meta["iteration"] = serde_json::json!(iter);
call_start_meta["llm_attempt"] = serde_json::json!(attempt);
}
if let Some(b) = bridge {
b.send_call_start(&call_id, "llm", "llm_call", call_start_meta);
}
dump_llm_request(
iteration.unwrap_or(0),
&call_id,
&effective_tool_format,
opts,
);
let first_token = super::first_token::FirstTokenTimer::for_current_span();
let start = std::time::Instant::now();
let detector_ctx = streaming_detector
.as_ref()
.map(|c| StreamingDetectorContext {
session_id: c.session_id.clone(),
known_tools: c.known_tools.clone(),
});
let raw_capture_context =
RawProviderCaptureContext::new(call_id.clone(), iteration.unwrap_or(0));
let llm_result = with_raw_provider_capture_context(raw_capture_context, async {
if let Some(b) = bridge {
let delta_tx = spawn_progress_forwarder(
b,
call_id.clone(),
user_visible,
detector_ctx,
first_token,
);
let delta_tx = match delta_sink.clone() {
Some(sink) => tee_delta_sender(vec![delta_tx, sink]),
None => delta_tx,
};
if offthread {
vm_call_llm_full_streaming_offthread_single_route(opts, delta_tx).await
} else {
vm_call_llm_full_streaming_single_route(opts, delta_tx).await
}
} else if offthread {
let delta_tx = match detector_ctx {
Some(ctx) => {
let detector_tx = spawn_detector_only_forwarder(ctx, first_token);
match delta_sink.clone() {
Some(sink) => tee_delta_sender(vec![detector_tx, sink]),
None => detector_tx,
}
}
None if let Some(sink) = delta_sink.clone() => sink,
None => {
let (tx, _rx) = tokio::sync::mpsc::unbounded_channel::<String>();
tx
}
};
vm_call_llm_full_streaming_offthread_single_route(opts, delta_tx).await
} else if let Some(sink) = delta_sink.clone() {
let delta_tx = match detector_ctx {
Some(ctx) => tee_delta_sender(vec![
spawn_detector_only_forwarder(ctx, first_token),
sink,
]),
None => sink,
};
vm_call_llm_full_streaming_single_route(opts, delta_tx).await
} else if let Some(ctx) = detector_ctx {
let delta_tx = spawn_detector_only_forwarder(ctx, first_token);
vm_call_llm_full_streaming_single_route(opts, delta_tx).await
} else {
vm_call_llm_full_single_route(opts).await
}
})
.await;
drop(rate_limit_permit);
let duration_ms = start.elapsed().as_millis() as u64;
record_governor_call_outcome(
&opts.provider,
&governor_org_key,
governor_reserved,
&llm_result,
);
match llm_result {
Ok(result) => {
if is_retryable_unproductive_completion(&result)
&& attempt < empty_completion_retry_budget(&opts.provider)
{
let errored_actionless = is_errored_actionless_completion(&result);
annotate_current_span(&[
("status", serde_json::json!("retrying")),
("retry_reason", serde_json::json!("empty_completion")),
("attempt", serde_json::json!(attempt)),
]);
let detail = if errored_actionless {
format!(
"provider {} model {} ended with a provider error (stop_reason=error) and emitted no tool call (the intended action went only to the reasoning channel)",
opts.provider, opts.model
)
} else {
format!(
"provider {} model {} returned a zero-token empty completion (no content, thinking, or tool calls)",
opts.provider, opts.model
)
};
let retry_reason = if errored_actionless {
UnproductiveCompletionReason::UnproductiveCompletion
} else {
UnproductiveCompletionReason::EmptyGeneration
};
emit_empty_completion_retry(
iteration.unwrap_or(0),
attempt + 1,
opts,
retry_reason,
duration_ms,
&detail,
);
if let Some(b) = bridge {
b.send_call_end(
&call_id,
"llm",
"llm_call",
duration_ms,
"retrying",
serde_json::json!({
"error": detail,
"retryable": true,
"attempt": attempt,
"user_visible": user_visible,
}),
);
}
attempt += 1;
empty_completion_retries += 1;
let backoff =
if crate::llm::providers::MockProvider::should_intercept(&opts.provider) {
0
} else {
base_retry_backoff_ms(attempt)
};
crate::events::log_warn(
"llm",
&format!("{detail}; retrying in {backoff}ms (attempt {attempt})"),
);
if backoff > 0 {
crate::clock_mock::sleep(std::time::Duration::from_millis(backoff)).await;
}
continue;
}
let attempt_count = attempt + 1;
let provider_under_throttle = provider_was_throttled_during_call
|| crate::llm::rate_governor::provider_already_throttled(
&opts.provider,
&governor_org_key,
);
if let Some(error) = terminal_unproductive_completion_failure(
opts,
&result,
provider_under_throttle,
attempt_count,
duration_ms,
) {
let category = crate::value::error_to_category(&error);
let message = error.to_string();
let classified = super::api::classify_llm_error(category.clone(), &message);
let status = "retries_exhausted";
annotate_current_span(&[
("status", serde_json::json!(status)),
("error", serde_json::json!(message.as_str())),
("retryable", serde_json::json!(false)),
("failover_eligible", serde_json::json!(true)),
("attempt", serde_json::json!(attempt)),
]);
dump_llm_response(
iteration.unwrap_or(0),
&call_id,
&result,
duration_ms,
opts.applied_structural_experiment.as_ref(),
opts.tools.as_ref(),
);
append_provider_call_error_observability(ProviderCallErrorObservation {
iteration: iteration.unwrap_or(0),
call_id: &call_id,
attempt,
status,
opts,
category: &category,
classified: &classified,
message: &message,
retryable: false,
failover_eligible: true,
attempt_count: Some(attempt_count),
});
dump_resolved_dispatch(
iteration.unwrap_or(0),
&call_id,
opts,
&effective_tool_format,
&super::resolved_dispatch::DispatchOutcome::from_error_message(&message),
);
if let Some(b) = bridge {
b.send_call_end(
&call_id,
"llm",
"llm_call",
duration_ms,
status,
serde_json::json!({
"error": message,
"retryable": false,
"failover_eligible": true,
"attempt": attempt,
"user_visible": user_visible,
}),
);
}
if let Some(metrics) = crate::active_metrics_registry() {
metrics.record_llm_call(
&result.provider,
&result.model,
status,
super::cost::calculate_cost_for_provider(
&result.provider,
&result.model,
result.input_tokens,
result.output_tokens,
),
);
}
return Err(error);
}
let usage = crate::tracing::LlmCallUsage {
model: result.model.clone(),
provider: result.provider.clone(),
input_tokens: result.input_tokens,
output_tokens: result.output_tokens,
cache_read_tokens: result.cache_read_tokens,
cache_write_tokens: result.cache_write_tokens,
cost_usd: result.priced_cost_usd(),
};
annotate_current_span(&[("status", serde_json::json!("ok"))]);
annotate_current_span(&usage.metadata_pairs());
dump_llm_response(
iteration.unwrap_or(0),
&call_id,
&result,
duration_ms,
opts.applied_structural_experiment.as_ref(),
opts.tools.as_ref(),
);
dump_resolved_dispatch(
iteration.unwrap_or(0),
&call_id,
opts,
&effective_tool_format,
&super::resolved_dispatch::DispatchOutcome::from_result(
&result,
empty_completion_retries,
),
);
annotate_current_span(&[(
"structural_experiment",
opts.applied_structural_experiment
.as_ref()
.map(serde_json::to_value)
.transpose()
.unwrap_or(None)
.unwrap_or(serde_json::Value::Null),
)]);
if let Some(b) = bridge {
b.send_call_end(
&call_id,
"llm",
"llm_call",
duration_ms,
"ok",
serde_json::json!({
"model": result.model,
"input_tokens": result.input_tokens,
"output_tokens": result.output_tokens,
"user_visible": user_visible,
"structural_experiment": opts.applied_structural_experiment.as_ref(),
}),
);
}
trace_llm_call(LlmTraceEntry {
model: result.model.clone(),
input_tokens: result.input_tokens,
output_tokens: result.output_tokens,
duration_ms,
});
if let Some(metrics) = crate::active_metrics_registry() {
metrics.record_llm_call(
&result.provider,
&result.model,
"succeeded",
super::cost::calculate_cost_for_provider(
&result.provider,
&result.model,
result.input_tokens,
result.output_tokens,
),
);
if result.cache_read_tokens > 0 {
metrics.record_llm_cache_hit(&result.provider);
}
}
super::trace::emit_agent_event(super::trace::AgentTraceEvent::LlmCall {
call_id: call_id.clone(),
model: result.model.clone(),
input_tokens: result.input_tokens,
output_tokens: result.output_tokens,
cache_tokens: result.cache_read_tokens,
duration_ms,
iteration: iteration.unwrap_or(0),
});
if is_retryable_unproductive_completion(&result)
&& !crate::llm::providers::is_internal_simulator(&opts.provider)
{
let reason = if is_empty_unproductive_completion(&result) {
UnproductiveCompletionReason::EmptyGeneration
} else {
UnproductiveCompletionReason::UnproductiveCompletion
};
super::rate_limit::observe_unproductive_completion_for_llm_call(
opts,
reason.as_str(),
);
} else {
super::rate_limit::observe_network_outcome_for_llm_call(opts, false);
}
return Ok(result);
}
Err(error) => {
let category = crate::value::error_to_category(&error);
let message = error.to_string();
let classified = super::api::classify_llm_error(category.clone(), &message);
super::rate_limit::observe_retry_after_for_llm_call(
opts,
shared_cooldown_ms_for_llm_error(&error),
);
let empty_completion_reason = empty_completion_retry_reason(&error);
let empty_completion_error = empty_completion_reason.is_some();
if !empty_completion_error {
super::rate_limit::observe_network_outcome_for_llm_call(
opts,
is_network_failure_llm_error(&error) || is_overloaded_llm_error(&error),
);
}
let retryable = is_retryable_llm_error(&error);
let empty_completion_retry = empty_completion_error
&& attempt < empty_completion_retry_budget(&opts.provider);
let native_tool_channel_degrade = !degraded_to_text
&& crate::llm_config::tool_format_channel(&effective_tool_format)
== Some(crate::llm_config::ToolFormatChannel::Native)
&& opts.native_tools.is_some()
&& (is_native_tool_channel_failure(&error)
|| is_billed_noncommittal_throw(&error));
let stream_transport_degrade = !degraded_stream_transport
&& !native_tool_channel_degrade
&& is_stream_transport_failure(&error)
&& can_degrade_stream_transport(opts);
let can_retry = empty_completion_retry
|| native_tool_channel_degrade
|| stream_transport_degrade;
let status = if can_retry {
"retrying"
} else if retryable {
"retries_exhausted"
} else {
"error"
};
annotate_current_span(&[
("status", serde_json::json!(status)),
("error", serde_json::json!(message.as_str())),
("retryable", serde_json::json!(retryable)),
("attempt", serde_json::json!(attempt)),
]);
append_provider_call_error_observability(ProviderCallErrorObservation {
iteration: iteration.unwrap_or(0),
call_id: &call_id,
attempt,
status,
opts,
category: &category,
classified: &classified,
message: &message,
retryable,
failover_eligible: false,
attempt_count: None,
});
if let Some(b) = bridge {
b.send_call_end(
&call_id,
"llm",
"llm_call",
duration_ms,
status,
serde_json::json!({
"error": error.to_string(),
"retryable": retryable,
"attempt": attempt,
"user_visible": user_visible,
}),
);
}
if !can_retry {
let surfaced_error = if empty_completion_error
&& !crate::llm::providers::is_internal_simulator(&opts.provider)
{
Some(provider_exhausted_error(
opts,
empty_completion_reason.expect("empty reason accompanies empty error"),
attempt + 1,
Some(duration_ms),
message.clone(),
))
} else {
None
};
if empty_completion_error
&& !crate::llm::providers::is_internal_simulator(&opts.provider)
{
super::rate_limit::observe_unproductive_completion_for_llm_call(
opts,
empty_completion_reason
.expect("terminal empty has a classified reason")
.as_str(),
);
}
if let Some(metrics) = crate::active_metrics_registry() {
metrics.record_llm_call(&opts.provider, &opts.model, status, 0.0);
}
dump_resolved_dispatch(
iteration.unwrap_or(0),
&call_id,
opts,
&effective_tool_format,
&super::resolved_dispatch::DispatchOutcome::from_error_message(&message),
);
return Err(surfaced_error.unwrap_or(error));
}
if empty_completion_error {
empty_completion_retries += 1;
emit_empty_completion_retry(
iteration.unwrap_or(0),
attempt + 1,
opts,
empty_completion_reason.expect("empty retry has a classified reason"),
duration_ms,
&error.to_string(),
);
}
let degraded_options =
native_tool_channel_degrade.then(|| degrade_options_to_text_channel(opts));
let stream_degraded_options = stream_transport_degrade
.then(|| degrade_options_to_non_streaming_transport(opts));
attempt += 1;
let backoff = llm_retry_backoff_ms(&error, attempt, &opts.provider);
crate::events::log_warn(
"llm",
&format!(
"LLM call failed ({error}), retrying in {backoff}ms (attempt {attempt})"
),
);
if let Some(degraded) = degraded_options {
let detail = format!(
"provider {} model {} native tool channel failed (server-side tool-call \
parser 500/EOF: {error}); degrading tool_format native -> json and \
retrying on the text channel",
degraded.provider, degraded.model
);
crate::events::log_warn("llm", &detail);
append_llm_observability_entry(
"tool_format_degrade",
serde_json::Map::from_iter([
(
"iteration".to_string(),
serde_json::json!(iteration.unwrap_or(0)),
),
("attempt".to_string(), serde_json::json!(attempt)),
("provider".to_string(), serde_json::json!(degraded.provider)),
("model".to_string(), serde_json::json!(degraded.model)),
("from".to_string(), serde_json::json!("native")),
("to".to_string(), serde_json::json!("json")),
("error".to_string(), serde_json::json!(error.to_string())),
]),
);
annotate_current_span(&[
("tool_format_degrade", serde_json::json!(true)),
("tool_format_degrade_from", serde_json::json!("native")),
("tool_format_degrade_to", serde_json::json!("json")),
]);
effective_tool_format = "json".to_string();
degraded_to_text = true;
working = std::borrow::Cow::Owned(degraded);
}
if let Some(degraded) = stream_degraded_options {
let detail = format!(
"provider {} model {} streaming transport failed ({error}); degrading \
stream=true -> false and retrying through request/response transport",
degraded.provider, degraded.model
);
crate::events::log_warn("llm", &detail);
append_llm_observability_entry(
"stream_transport_degrade",
serde_json::Map::from_iter([
(
"iteration".to_string(),
serde_json::json!(iteration.unwrap_or(0)),
),
("attempt".to_string(), serde_json::json!(attempt)),
("provider".to_string(), serde_json::json!(degraded.provider)),
("model".to_string(), serde_json::json!(degraded.model)),
("from".to_string(), serde_json::json!(true)),
("to".to_string(), serde_json::json!(false)),
("error".to_string(), serde_json::json!(error.to_string())),
]),
);
annotate_current_span(&[
("stream_transport_degrade", serde_json::json!(true)),
("stream_transport_degrade_from", serde_json::json!(true)),
("stream_transport_degrade_to", serde_json::json!(false)),
]);
degraded_stream_transport = true;
working = std::borrow::Cow::Owned(degraded);
}
if backoff > 0 {
crate::clock_mock::sleep(std::time::Duration::from_millis(backoff)).await;
}
}
}
}
}
#[cfg(test)]
#[path = "agent_observe_cost_tests.rs"]
mod cost_tests;
#[cfg(test)]
mod empty_completion_retry_tests;
#[cfg(test)]
mod retry_tests;
#[cfg(test)]
mod streaming_detector_tests;