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
61 changes: 56 additions & 5 deletions crates/goose/src/agents/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ use crate::session::{Session, SessionManager, SessionNameUpdate};
use crate::tool_inspection::ToolInspectionManager;
use crate::tool_monitor::RepetitionInspector;
use crate::utils::is_token_cancelled;
use goose_providers::conversation::token_usage::ProviderUsage;
use goose_providers::conversation::token_usage::{ProviderUsage, Usage};
use goose_providers::errors::ProviderError;
use goose_providers::thinking::ThinkingEffort;
use regex::Regex;
Expand Down Expand Up @@ -1623,7 +1623,16 @@ impl Agent {

#[instrument(
skip(self, user_message, session_config, cancel_token),
fields(user_message, trace_input, session.id = %session_config.id)
fields(
user_message,
trace_input,
session.id = %session_config.id,
gen_ai.operation.name = "invoke_agent",
gen_ai.input.messages = tracing::field::Empty,
gen_ai.output.messages = tracing::field::Empty,
gen_ai.usage.input_tokens = tracing::field::Empty,
gen_ai.usage.output_tokens = tracing::field::Empty,
)
)]
pub async fn reply(
&self,
Expand Down Expand Up @@ -1652,6 +1661,12 @@ impl Agent {
let message_text_for_trace = agent_visible_message_text(&user_message);
tracing::Span::current().record("user_message", message_text_for_trace.as_str());
tracing::Span::current().record("trace_input", message_text_for_trace.as_str());
if gen_ai_telemetry::capture_message_content() {
tracing::Span::current().record(
"gen_ai.input.messages",
gen_ai_telemetry::simple_input_json(&message_text_for_trace).as_str(),
);
}

for content in &user_message.content {
if let MessageContent::ActionRequired(action_required) = content {
Expand Down Expand Up @@ -1865,6 +1880,7 @@ impl Agent {
.await?;

let conversation_to_compact = conversation.clone();
let reply_span = tracing::Span::current();

Ok(Box::pin(async_stream::try_stream! {
for event in command_preamble {
Expand Down Expand Up @@ -1936,7 +1952,7 @@ impl Agent {
}
};

let mut reply_stream = self.reply_internal(final_conversation, session_config, session, cancel_token).await?;
let mut reply_stream = self.reply_internal(final_conversation, session_config, session, cancel_token, reply_span.clone()).await?;
while let Some(event) = reply_stream.next().await {
yield event?;
}
Expand All @@ -1949,6 +1965,7 @@ impl Agent {
session_config: SessionConfig,
session: Session,
cancel_token: Option<CancellationToken>,
reply_span: tracing::Span,
) -> Result<BoxStream<'_, Result<AgentEvent>>> {
let context = self
.prepare_reply_context(&session.id, conversation, session.working_dir.as_path())
Expand Down Expand Up @@ -1979,7 +1996,7 @@ impl Agent {
.ok()
.and_then(|model_info| model_info.resolved_model)
.map(|resolved_model| InferenceMetadata {
provider: provider_name,
provider: provider_name.clone(),
requested_model,
resolved_model: Some(resolved_model),
});
Expand Down Expand Up @@ -2018,14 +2035,35 @@ impl Agent {

let working_dir = session.working_dir.clone();
let reply_stream_span = tracing::info_span!(
target: "goose::agents::agent",
parent: &reply_span,
"reply_stream",
trace_output = tracing::field::Empty,
session.id = %session_config.id,
session.user = %crate::session_context::session_user(),
session.host = %crate::session_context::session_host(),
session.agent_type = "goose",
gen_ai.operation.name = "invoke_agent",
gen_ai.conversation.id = %session_config.id,
gen_ai.request.model = %model_config.model_name,
gen_ai.provider.name = %provider_name,
gen_ai.input.messages = tracing::field::Empty,
gen_ai.output.messages = tracing::field::Empty,
gen_ai.usage.input_tokens = tracing::field::Empty,
gen_ai.usage.output_tokens = tracing::field::Empty,
Comment on lines +2051 to +2052

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Declare cache-token fields on the aggregate spans

When a turn includes cache read/write tokens, the new aggregate reply/reply_stream usage recording drops them: record_usage records gen_ai.usage.cache_read.input_tokens and gen_ai.usage.cache_creation.input_tokens, but tracing only accepts fields declared at the span callsite. This span (and the reply span above) declares only input/output tokens, so cache-token records are ignored and the root telemetry remains incomplete for providers that report cache usage.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Good feature for a future PR

);
if gen_ai_telemetry::capture_message_content() {
if let Some(last_user_msg) = conversation
.messages()
.iter()
.rev()
.find(|m| m.role == rmcp::model::Role::User)
{
reply_stream_span.record(
"gen_ai.input.messages",
gen_ai_telemetry::simple_input_json(&last_user_msg.as_concat_text()).as_str(),
);
}
}
let inner = Box::pin(async_stream::try_stream! {
let mut turns_taken = 0u32;
let max_turns = session_config.max_turns.unwrap_or_else(|| {
Expand All @@ -2037,6 +2075,7 @@ impl Agent {
let mut empty_turn_retries = 0u32;
let mut retrying_after_empty_turn = false;
let mut last_assistant_text = String::new();
let mut turn_total_usage = Usage::default();
let mut goal_check_pending = false;
let mut tool_pair_summarization_done = false;
let mut stop_hook_handled_for_exit = false;
Expand Down Expand Up @@ -2212,6 +2251,7 @@ impl Agent {
if let Some(ref usage) = usage {
let enriched = self.update_session_metrics(&session_config.id, session_config.schedule_id.clone(), usage, None).await?;
yield AgentEvent::Usage(enriched.clone());
turn_total_usage += enriched.usage;
pending_turn_usage = Some(enriched);
}

Expand Down Expand Up @@ -3083,7 +3123,18 @@ impl Agent {

if !last_assistant_text.is_empty() {
tracing::Span::current().record("trace_output", last_assistant_text.as_str());
if gen_ai_telemetry::capture_message_content() {
let output_json =
gen_ai_telemetry::simple_output_json(&last_assistant_text);
tracing::Span::current().record(
"gen_ai.output.messages",
output_json.as_str(),
);
reply_span.record("gen_ai.output.messages", output_json.as_str());
}
}
gen_ai_telemetry::record_usage(&tracing::Span::current(), &turn_total_usage);
gen_ai_telemetry::record_usage(&reply_span, &turn_total_usage);

if !stop_hook_handled_for_exit {
self.emit_stop_hook(&session_config.id, &last_assistant_text, &session.working_dir.to_string_lossy()).await;
Expand Down
53 changes: 45 additions & 8 deletions crates/goose/src/agents/gen_ai_telemetry.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
use crate::conversation::message::{Message, MessageContent, ToolResult};
use goose_providers::conversation::token_usage::ProviderUsage;
use goose_providers::conversation::token_usage::{ProviderUsage, Usage};
use rmcp::model::{CallToolRequestParams, CallToolResult, Role};
use serde_json::{json, Value};
use tracing::Span;
Expand All @@ -15,6 +15,14 @@ pub(super) fn input_messages_json(messages: &[Message]) -> String {
Value::Array(messages.iter().map(message_json).collect()).to_string()
}

pub(super) fn simple_input_json(text: &str) -> String {
json!([{"role": "user", "content": text}]).to_string()
}

pub(super) fn simple_output_json(text: &str) -> String {
json!([{"role": "assistant", "content": text, "finish_reason": "stop"}]).to_string()
}

pub(super) fn output_message_json(message: &Message) -> String {
// Message does not retain provider finish reasons; tool requests are the only
// distinct completion signal available after streaming.
Expand All @@ -40,22 +48,26 @@ pub(super) fn append_message(accumulated: &mut Option<Message>, message: &Messag
}
}

pub(super) fn record_provider_usage(span: &Span, usage: &ProviderUsage) {
span.record("gen_ai.response.model", usage.model.as_str());
if let Some(tokens) = usage.usage.input_tokens {
pub(super) fn record_usage(span: &Span, usage: &Usage) {
if let Some(tokens) = usage.input_tokens {
span.record("gen_ai.usage.input_tokens", tokens);
}
if let Some(tokens) = usage.usage.output_tokens {
if let Some(tokens) = usage.output_tokens {
span.record("gen_ai.usage.output_tokens", tokens);
}
if let Some(tokens) = usage.usage.cache_read_input_tokens {
if let Some(tokens) = usage.cache_read_input_tokens {
span.record("gen_ai.usage.cache_read.input_tokens", tokens);
}
if let Some(tokens) = usage.usage.cache_write_input_tokens {
if let Some(tokens) = usage.cache_write_input_tokens {
span.record("gen_ai.usage.cache_creation.input_tokens", tokens);
}
}

pub(super) fn record_provider_usage(span: &Span, usage: &ProviderUsage) {
span.record("gen_ai.response.model", usage.model.as_str());
record_usage(span, &usage.usage);
}

pub(super) fn tool_result_json(result: &ToolResult<CallToolResult>) -> String {
match result {
Ok(result) if result.is_error != Some(true) => json!({
Expand Down Expand Up @@ -98,13 +110,38 @@ fn message_json(message: &Message) -> Value {
}
};

let parts: Vec<Value> = message.content.iter().map(message_part_json).collect();
let parts = consolidated_parts(&message.content);
json!({
"role": role,
"parts": parts,
})
}

/// Merge consecutive text and reasoning parts into single entries so that
/// streaming tokens don't each get their own JSON object in the OTEL output.
fn consolidated_parts(content: &[MessageContent]) -> Vec<Value> {
let mut result: Vec<Value> = Vec::new();
for item in content {
let value = message_part_json(item);
let item_type = value.get("type").and_then(|v| v.as_str());
if matches!(item_type, Some("text" | "reasoning")) {
if let Some(last) = result.last_mut() {
if last.get("type") == value.get("type") {
if let (Some(existing), Some(new_content)) = (
last.get("content").and_then(|v| v.as_str()),
value.get("content").and_then(|v| v.as_str()),
) {
last["content"] = Value::String(format!("{}{}", existing, new_content));
continue;
}
}
}
}
result.push(value);
}
result
}

fn tool_call_part(id: &str, tool_call: &ToolResult<CallToolRequestParams>) -> Value {
match tool_call {
Ok(tool_call) => json!({
Expand Down
3 changes: 1 addition & 2 deletions crates/goose/src/agents/reply_parts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -879,8 +879,7 @@ mod tests {
let output: Value =
serde_json::from_str(fields["gen_ai.output.messages"].as_str().unwrap()).unwrap();
assert_eq!(output[0]["finish_reason"], "stop");
assert_eq!(output[0]["parts"][0]["content"], "hello ");
assert_eq!(output[0]["parts"][1]["content"], "world");
assert_eq!(output[0]["parts"][0]["content"], "hello world");
}

#[tokio::test]
Expand Down