Skip to main content

aether_telemetry/
otel_observer.rs

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