From 5d0cf5b070833c7fa1d834ae679caee769a7ff99 Mon Sep 17 00:00:00 2001 From: xiaoxixi Date: Fri, 24 Jul 2026 17:21:44 +0800 Subject: [PATCH] feat(agent): record tool/turn/provider metrics --- src/agent/agent_loop.rs | 65 +++++++++++++++++++++++++++++++++++------ 1 file changed, 56 insertions(+), 9 deletions(-) diff --git a/src/agent/agent_loop.rs b/src/agent/agent_loop.rs index a70c1dd..cc4f2bf 100644 --- a/src/agent/agent_loop.rs +++ b/src/agent/agent_loop.rs @@ -491,16 +491,41 @@ impl AgentLoop { iteration: u32, turn: Option<&AgentTurnContext>, ) -> Result { - let mut provider_stream = self.provider.stream(request).await.map_err(|error| { - tracing::error!(error = %error, "LLM request failed"); - AgentError::LlmError(error.to_string()) - })?; + let metrics = crate::observability::metrics::global_metrics(); + let provider_name = self.provider.name().to_string(); + let provider_model = self.provider.model_id().to_string(); + let start = Instant::now(); + + let mut provider_stream = match self.provider.stream(request).await { + Ok(stream) => stream, + Err(error) => { + tracing::error!(error = %error, "LLM request failed"); + metrics.record_provider( + &provider_name, + &provider_model, + None, + start.elapsed().as_millis() as u64, + true, + ); + return Err(AgentError::LlmError(error.to_string())); + } + }; let mut accumulator = ProviderResponseAccumulator::default(); while let Some(chunk) = provider_stream.next().await { - let chunk = chunk.map_err(|error| { - tracing::error!(error = %error, "LLM stream failed"); - AgentError::LlmError(error.to_string()) - })?; + let chunk = match chunk { + Ok(chunk) => chunk, + Err(error) => { + tracing::error!(error = %error, "LLM stream failed"); + metrics.record_provider( + &provider_name, + &provider_model, + None, + start.elapsed().as_millis() as u64, + true, + ); + return Err(AgentError::LlmError(error.to_string())); + } + }; if let Some(turn) = turn { let event = match &chunk { ProviderChunk::Reasoning(delta) => Some(TurnEvent::ReasoningDelta { @@ -521,7 +546,11 @@ impl AgentLoop { } accumulator.push(chunk); } - Ok(accumulator.finish()) + let response = accumulator.finish(); + let latency_ms = start.elapsed().as_millis() as u64; + metrics.record_provider(&provider_name, &provider_model, None, latency_ms, false); + metrics.record_provider_tokens(&provider_name, &response.usage); + Ok(response) } fn annotate_message( @@ -615,6 +644,8 @@ impl AgentLoop { mut messages: Vec, turn: Option, ) -> Result { + let turn_start = Instant::now(); + #[cfg(debug_assertions)] tracing::debug!( history_len = messages.len(), @@ -703,6 +734,10 @@ impl AgentLoop { assistant_message.provider_state = response.provider_state; Self::annotate_message(&mut assistant_message, turn.as_ref(), iteration, true); emitted_messages.push(assistant_message.clone()); + crate::observability::metrics::global_metrics().record_turn( + Some(&accumulated_usage), + turn_start.elapsed().as_millis() as u64, + ); return Ok(AgentProcessResult { final_response: assistant_message, emitted_messages, @@ -844,6 +879,10 @@ impl AgentLoop { true, ); emitted_messages.push(assistant_message.clone()); + crate::observability::metrics::global_metrics().record_turn( + Some(&accumulated_usage), + turn_start.elapsed().as_millis() as u64, + ); Ok(AgentProcessResult { final_response: assistant_message, emitted_messages, @@ -876,6 +915,11 @@ impl AgentLoop { let mut final_message = ChatMessage::assistant(fallback); Self::annotate_message(&mut final_message, turn.as_ref(), summary_iteration, true); emitted_messages.push(final_message.clone()); + let turn_usage = (accumulated_usage.total_tokens > 0).then_some(&accumulated_usage); + crate::observability::metrics::global_metrics().record_turn( + turn_usage, + turn_start.elapsed().as_millis() as u64, + ); Ok(AgentProcessResult { final_response: final_message, emitted_messages, @@ -1016,6 +1060,9 @@ impl AgentLoop { }); } + crate::observability::metrics::global_metrics() + .record_tool_call(&tool_name, result.success); + // Apply duration Ok(ToolExecutionOutcome { duration, ..result }) }