Skip to main content

aether_telemetry/
otel_observer.rs

1use crate::content_capture::{ContentBuffer, ContentCapture};
2use crate::content_json::{input_messages_json, output_messages_json, tool_definitions_json};
3use crate::gen_ai_metrics::GenAiMetrics;
4use crate::genai_constants as semconv;
5use crate::llm_call_state::LlmCallState;
6use crate::span_guard::{ErrorKind, SpanGuard};
7use crate::trace_context::inject_trace_context;
8use aether_core::events::{
9    AgentEvent, AgentObserver, LlmCallPurpose, MessageEvent, ToolEvent, TraceContext, TurnEvent, TurnOutcome,
10    mcp_tool_name,
11};
12use llm::catalog::Provider;
13use llm::{ContentBlock, ModelPricing, ToolCallError, ToolCallRequest, ToolCallResult, ToolDefinition};
14use opentelemetry::trace::{SpanBuilder, SpanKind, TraceContextExt, Tracer as _};
15use opentelemetry::{Context, KeyValue};
16use opentelemetry_sdk::trace::SdkTracer;
17use std::collections::HashMap;
18
19/// [`AgentObserver`] that renders an agent's event stream as OpenTelemetry
20/// `GenAI` spans and metrics.
21pub struct OtelObserver {
22    turn: Option<TurnState>,
23    tool_definitions: Vec<ToolDefinition>,
24    otel: OtelInstrumentation,
25}
26
27/// The process-wide export pipeline — tracer, meters, and the capture and
28/// parenting policy every span it produces inherits.
29#[derive(Clone)]
30pub struct OtelInstrumentation {
31    pub tracer: SdkTracer,
32    pub metrics: GenAiMetrics,
33    pub capture_content: bool,
34    pub root_parent: Option<Context>,
35    pub agent_name: Option<String>,
36}
37
38impl OtelObserver {
39    pub fn new(otel: OtelInstrumentation) -> Self {
40        Self { turn: None, tool_definitions: Vec::new(), otel }
41    }
42}
43
44impl AgentObserver for OtelObserver {
45    fn on_event(&mut self, message: &AgentEvent) {
46        match message {
47            AgentEvent::Turn(TurnEvent::Started { content }) => self.start_turn(content),
48            AgentEvent::Turn(TurnEvent::Ended { outcome }) => {
49                if let Some(turn) = self.turn.take() {
50                    turn.finish(outcome);
51                }
52            }
53            AgentEvent::Tool(ToolEvent::DefinitionsUpdated { tools }) => {
54                self.tool_definitions.clone_from(tools);
55            }
56            message => {
57                if let Some(turn) = &mut self.turn {
58                    turn.on_event(message, &self.otel, &self.tool_definitions);
59                }
60            }
61        }
62    }
63
64    fn tool_trace_context(&self, tool_id: &str) -> Option<TraceContext> {
65        let span = self.turn.as_ref()?.executing_tools.get(tool_id)?;
66        inject_trace_context(span.context())
67    }
68}
69
70impl OtelObserver {
71    fn start_turn(&mut self, content: &[ContentBlock]) {
72        // Drop (and thereby cancel) any stale turn and its open spans before
73        // starting the new span, so the stale turn can't become its parent.
74        self.turn = None;
75
76        let mut input = self.otel.capture().buffer();
77        input.set(&ContentBlock::join_text(content));
78        let operation_name = "invoke_agent";
79        let mut attributes = vec![KeyValue::new(semconv::GEN_AI_OPERATION_NAME, operation_name)];
80        let span_name = match &self.otel.agent_name {
81            Some(agent_name) => {
82                attributes.push(KeyValue::new(semconv::GEN_AI_AGENT_NAME, agent_name.clone()));
83                format!("{operation_name} {agent_name}")
84            }
85            None => operation_name.to_string(),
86        };
87        if let Some(text) = input.get() {
88            attributes.push(KeyValue::new(semconv::GEN_AI_INPUT_MESSAGES, input_messages_json(text)));
89        }
90        let builder = SpanBuilder::from_name(span_name).with_kind(SpanKind::Internal).with_attributes(attributes);
91        let span_context = self.otel.start_span(builder, self.otel.root_parent.as_ref());
92        let span = SpanGuard::new(span_context, TURN_CANCEL_MESSAGE);
93        self.turn = Some(TurnState::new(span, input, self.otel.capture()));
94    }
95}
96
97impl OtelInstrumentation {
98    fn capture(&self) -> ContentCapture {
99        ContentCapture::from_enabled(self.capture_content)
100    }
101
102    /// Starts a span under `parent`. Providers for disabled tracing use an
103    /// always-off sampler, so callers never need to branch.
104    pub(crate) fn start_span(&self, builder: SpanBuilder, parent: Option<&Context>) -> Context {
105        let parent = parent.cloned().unwrap_or_default();
106        Context::new().with_span(self.tracer.build_with_context(builder, &parent))
107    }
108}
109
110const TURN_CANCEL_MESSAGE: &str = "turn cancelled";
111const TOOL_CANCEL_MESSAGE: &str = "turn ended before the tool completed";
112
113/// All state scoped to one turn: the turn span plus any in-flight LLM-call and
114/// tool spans and captured content. A turn exists by construction inside its
115/// methods, and dropping it ends every still-open span as cancelled, so
116/// replacing the value is all a reset takes.
117struct TurnState {
118    span: SpanGuard,
119    input: ContentBuffer,
120    output: ContentBuffer,
121    chat_call: Option<LlmCallState>,
122    compaction_call: Option<LlmCallState>,
123    /// Tool-call arguments streamed before execution starts, keyed by call id.
124    streamed_arguments: HashMap<String, ContentBuffer>,
125    /// Spans of currently executing tools, keyed by call id.
126    executing_tools: HashMap<String, SpanGuard>,
127}
128
129impl TurnState {
130    fn new(span: SpanGuard, input: ContentBuffer, capture: ContentCapture) -> Self {
131        Self {
132            span,
133            input,
134            output: capture.buffer(),
135            chat_call: None,
136            compaction_call: None,
137            streamed_arguments: HashMap::new(),
138            executing_tools: HashMap::new(),
139        }
140    }
141
142    fn on_event(&mut self, message: &AgentEvent, instrumentation: &OtelInstrumentation, tools: &[ToolDefinition]) {
143        match message {
144            AgentEvent::Turn(TurnEvent::LlmCallStarted {
145                purpose,
146                provider,
147                model,
148                display_name,
149                pricing,
150                attempt,
151                ..
152            }) => {
153                self.start_llm_call(
154                    LlmCallStart {
155                        purpose: *purpose,
156                        provider: provider.as_deref(),
157                        model: model.as_deref(),
158                        display_name,
159                        pricing: *pricing,
160                        attempt: *attempt,
161                    },
162                    instrumentation,
163                    tools,
164                );
165            }
166            AgentEvent::Turn(TurnEvent::LlmCallEnded { purpose, outcome }) => {
167                if let Some(call) = self.llm_call_slot(*purpose).take() {
168                    call.finish(outcome);
169                }
170            }
171            AgentEvent::Message(
172                MessageEvent::Text { message_id, chunk, is_complete: false }
173                | MessageEvent::Thought { message_id, chunk, is_complete: false },
174            ) => {
175                if let Some(chat) = &mut self.chat_call {
176                    chat.record_response_chunk(message_id, chunk);
177                }
178            }
179            AgentEvent::Message(MessageEvent::Text { chunk, is_complete: true, .. }) => {
180                self.output.push(chunk);
181            }
182            AgentEvent::Tool(ToolEvent::Call { request, .. }) => self.on_tool_call(request, instrumentation),
183            AgentEvent::Tool(ToolEvent::CallUpdate { tool_call_id, chunk, .. }) => {
184                self.on_tool_call_update(tool_call_id, chunk);
185            }
186            AgentEvent::Tool(ToolEvent::ExecutionStarted { tool_id, tool_name }) => {
187                self.on_tool_execution_started(tool_id, tool_name, instrumentation);
188            }
189            AgentEvent::Tool(ToolEvent::Result { result, .. }) => self.on_tool_result(result, instrumentation),
190            AgentEvent::Tool(ToolEvent::Error { error, .. }) => self.on_tool_error(error),
191            _ => {}
192        }
193    }
194
195    fn finish(self, outcome: &TurnOutcome) {
196        let Self { mut span, output, chat_call, compaction_call, executing_tools, .. } = self;
197        // Cancel any still-open call and tool spans before ending their parent.
198        drop(chat_call);
199        drop(compaction_call);
200        drop(executing_tools);
201        match outcome {
202            TurnOutcome::Completed => {
203                if let Some(text) = output_messages_json(output.get(), &[], None) {
204                    span.set_attribute(KeyValue::new(semconv::GEN_AI_OUTPUT_MESSAGES, text));
205                }
206                span.end_ok();
207            }
208            TurnOutcome::Failed { error } => span.end_error(None, error.clone()),
209            TurnOutcome::Cancelled => span.end_error(Some(ErrorKind::Cancelled), TURN_CANCEL_MESSAGE),
210        }
211    }
212
213    fn start_llm_call(
214        &mut self,
215        call: LlmCallStart<'_>,
216        instrumentation: &OtelInstrumentation,
217        tools: &[ToolDefinition],
218    ) {
219        let model_name = call.model.unwrap_or(call.display_name).to_string();
220
221        // Metrics backends key series by attribute values, so the metric set
222        // stays a small fixed vocabulary — no content, no per-call details.
223        let mut metric_attributes = vec![
224            KeyValue::new(semconv::GEN_AI_OPERATION_NAME, "chat"),
225            KeyValue::new(semconv::GEN_AI_REQUEST_MODEL, model_name.clone()),
226        ];
227
228        if let Some(provider) = call.provider {
229            metric_attributes.push(KeyValue::new(semconv::GEN_AI_PROVIDER_NAME, genai_provider_name(provider)));
230        }
231
232        if call.purpose == LlmCallPurpose::Compaction {
233            metric_attributes.push(KeyValue::new(semconv::LLM_PURPOSE, "compaction"));
234        }
235
236        let mut attributes = metric_attributes.clone();
237        attributes.push(KeyValue::new(semconv::GEN_AI_REQUEST_STREAM, true));
238        attributes.push(KeyValue::new(semconv::LLM_ATTEMPT, i64::from(call.attempt)));
239        if let Some(pricing) = call.pricing {
240            attributes.extend(pricing_attributes(pricing));
241        }
242        // Only chat calls carry the turn's input and tool definitions; a
243        // compaction call's actual input is the internal summarization prompt.
244        if call.purpose == LlmCallPurpose::Chat {
245            if let Some(input) = self.input.get() {
246                attributes.push(KeyValue::new(semconv::GEN_AI_INPUT_MESSAGES, input_messages_json(input)));
247            }
248            if instrumentation.capture_content && !tools.is_empty() {
249                attributes.push(KeyValue::new(semconv::GEN_AI_TOOL_DEFINITIONS, tool_definitions_json(tools)));
250            }
251        }
252
253        let name = if model_name.is_empty() { "chat".to_string() } else { format!("chat {model_name}") };
254        let builder = SpanBuilder::from_name(name).with_kind(SpanKind::Client).with_attributes(attributes);
255        let context = instrumentation.start_span(builder, Some(self.span.context()));
256        let state = LlmCallState::new(
257            context,
258            instrumentation.metrics.clone(),
259            call.purpose,
260            instrumentation.capture(),
261            metric_attributes,
262        );
263        *self.llm_call_slot(call.purpose) = Some(state);
264    }
265
266    fn on_tool_call(&mut self, request: &ToolCallRequest, instrumentation: &OtelInstrumentation) {
267        if let Some(chat) = &mut self.chat_call {
268            chat.record_tool_call_start(request);
269        }
270        let mut arguments = instrumentation.capture().buffer();
271        arguments.set(&request.arguments);
272        self.streamed_arguments.insert(request.id.clone(), arguments);
273    }
274
275    fn on_tool_call_update(&mut self, tool_call_id: &str, chunk: &str) {
276        if let Some(chat) = &mut self.chat_call {
277            chat.record_tool_call_update(tool_call_id, chunk);
278        }
279        if let Some(arguments) = self.streamed_arguments.get_mut(tool_call_id) {
280            arguments.push(chunk);
281        }
282    }
283
284    fn on_tool_execution_started(&mut self, tool_id: &str, tool_name: &str, instrumentation: &OtelInstrumentation) {
285        let mut attributes = vec![
286            KeyValue::new(semconv::GEN_AI_OPERATION_NAME, "execute_tool"),
287            KeyValue::new(semconv::GEN_AI_TOOL_NAME, tool_name.to_string()),
288            KeyValue::new(semconv::GEN_AI_TOOL_CALL_ID, tool_id.to_string()),
289            KeyValue::new(semconv::MCP_METHOD_NAME, "tools/call"),
290            // The server sees the tool without our namespacing prefix, so this
291            // is what joins this span to the server's span for the same call.
292            KeyValue::new(semconv::MCP_TOOL_NAME, mcp_tool_name(tool_name).to_string()),
293        ];
294        let arguments = self.streamed_arguments.remove(tool_id);
295
296        if let Some(text) = arguments.as_ref().and_then(ContentBuffer::get) {
297            attributes.push(KeyValue::new(semconv::GEN_AI_TOOL_CALL_ARGUMENTS, text.to_string()));
298        }
299
300        let builder = SpanBuilder::from_name(format!("execute_tool {tool_name}"))
301            .with_kind(SpanKind::Client)
302            .with_attributes(attributes);
303        let context = instrumentation.start_span(builder, Some(self.span.context()));
304        self.executing_tools.insert(tool_id.to_string(), SpanGuard::new(context, TOOL_CANCEL_MESSAGE));
305    }
306
307    fn on_tool_result(&mut self, result: &ToolCallResult, instrumentation: &OtelInstrumentation) {
308        self.streamed_arguments.remove(&result.id);
309        let Some(mut span) = self.executing_tools.remove(&result.id) else { return };
310
311        if let Some(content) = instrumentation.capture().content(&result.result) {
312            span.set_attribute(KeyValue::new(semconv::GEN_AI_TOOL_CALL_RESULT, content.to_string()));
313        }
314
315        span.end_ok();
316    }
317
318    fn on_tool_error(&mut self, error: &ToolCallError) {
319        self.streamed_arguments.remove(&error.id);
320        if let Some(mut span) = self.executing_tools.remove(&error.id) {
321            span.end_error(Some(ErrorKind::ToolError), error.error.clone());
322        }
323    }
324
325    fn llm_call_slot(&mut self, purpose: LlmCallPurpose) -> &mut Option<LlmCallState> {
326        match purpose {
327            LlmCallPurpose::Chat => &mut self.chat_call,
328            LlmCallPurpose::Compaction => &mut self.compaction_call,
329        }
330    }
331}
332
333/// Maps the agent's canonical provider name to the `GenAI` semantic-convention
334/// provider name, passing providers the catalog doesn't know through as-is.
335fn genai_provider_name(provider: &str) -> String {
336    provider.parse::<Provider>().map_or_else(|_| provider.to_string(), |p| p.genai_provider_name().to_string())
337}
338
339/// Borrowed view of [`TurnEvent::LlmCallStarted`]; named fields keep the two
340/// optional strings from being swapped at the call site.
341#[derive(Clone, Copy)]
342struct LlmCallStart<'a> {
343    purpose: LlmCallPurpose,
344    provider: Option<&'a str>,
345    model: Option<&'a str>,
346    display_name: &'a str,
347    pricing: Option<ModelPricing>,
348    attempt: u32,
349}
350
351fn pricing_attributes(pricing: ModelPricing) -> Vec<KeyValue> {
352    let mut attributes = vec![
353        KeyValue::new(semconv::AI_INPUT_TOKEN_PRICE, pricing.input_per_million / TOKENS_PER_MILLION),
354        KeyValue::new(semconv::AI_OUTPUT_TOKEN_PRICE, pricing.output_per_million / TOKENS_PER_MILLION),
355    ];
356
357    if let Some(price) = pricing.cache_read_per_million {
358        attributes.push(KeyValue::new(semconv::AI_CACHE_READ_TOKEN_PRICE, price / TOKENS_PER_MILLION));
359    }
360
361    if let Some(price) = pricing.cache_write_per_million {
362        attributes.push(KeyValue::new(semconv::AI_CACHE_WRITE_TOKEN_PRICE, price / TOKENS_PER_MILLION));
363    }
364
365    attributes
366}
367
368const TOKENS_PER_MILLION: f64 = 1_000_000.0;