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    /// Tool-call arguments streamed before execution starts, keyed by call id.
131    streamed_arguments: HashMap<String, ContentBuffer>,
132    /// Spans of currently executing tools, keyed by call id.
133    executing_tools: HashMap<String, SpanGuard>,
134}
135
136impl TurnState {
137    fn new(span: SpanGuard, input: ContentBuffer, capture_output: bool) -> Self {
138        Self {
139            span,
140            input,
141            output: ContentBuffer::new(capture_output),
142            chat_call: None,
143            compaction_call: None,
144            streamed_arguments: HashMap::new(),
145            executing_tools: HashMap::new(),
146        }
147    }
148
149    fn on_event(
150        &mut self,
151        message: &AgentEvent,
152        instrumentation: &OtelInstrumentation,
153        tools: &[ToolDefinition],
154        system_prompt: Option<&str>,
155    ) {
156        match message {
157            AgentEvent::Turn(TurnEvent::LlmCallStarted { purpose, model, display_name, attempt, .. }) => {
158                self.start_llm_call(
159                    LlmCallStart {
160                        purpose: *purpose,
161                        provider: model.provider.as_deref(),
162                        model: model.model_id.as_deref(),
163                        display_name,
164                        pricing: model.pricing,
165                        attempt: *attempt,
166                    },
167                    instrumentation,
168                    tools,
169                    system_prompt,
170                );
171            }
172            AgentEvent::Turn(TurnEvent::LlmCallEnded { purpose, outcome }) => {
173                if let Some(call) = self.llm_call_slot(*purpose).take() {
174                    call.finish(outcome);
175                }
176            }
177            AgentEvent::Message(
178                MessageEvent::Text { message_id, chunk, is_complete: false }
179                | MessageEvent::Thought { message_id, chunk, is_complete: false },
180            ) => {
181                if let Some(chat) = &mut self.chat_call {
182                    chat.record_response_chunk(message_id, chunk);
183                }
184            }
185            AgentEvent::Message(MessageEvent::Text { chunk, is_complete: true, .. }) => {
186                self.output.push(chunk);
187            }
188            AgentEvent::Tool(ToolEvent::Call { request, .. }) => self.on_tool_call(request, instrumentation),
189            AgentEvent::Tool(ToolEvent::CallUpdate { tool_call_id, chunk, .. }) => {
190                self.on_tool_call_update(tool_call_id, chunk);
191            }
192            AgentEvent::Tool(ToolEvent::ExecutionStarted { tool_id, tool_name }) => {
193                self.on_tool_execution_started(tool_id, tool_name, instrumentation);
194            }
195            AgentEvent::Tool(ToolEvent::Result { result, .. }) => self.on_tool_result(result, instrumentation),
196            AgentEvent::Tool(ToolEvent::Error { error, .. }) => self.on_tool_error(error),
197            _ => {}
198        }
199    }
200
201    fn finish(self, outcome: &TurnOutcome) {
202        let Self { mut span, output, chat_call, compaction_call, executing_tools, .. } = self;
203        // Cancel any still-open call and tool spans before ending their parent.
204        drop(chat_call);
205        drop(compaction_call);
206        drop(executing_tools);
207        match outcome {
208            TurnOutcome::Completed => {
209                if let Some(text) = output_messages_json(output.get(), &[], None) {
210                    span.set_attribute(KeyValue::new(semconv::GEN_AI_OUTPUT_MESSAGES, text));
211                }
212                span.end_ok();
213            }
214            TurnOutcome::Failed { error } => span.end_error(None, error.clone()),
215            TurnOutcome::Cancelled => span.end_error(Some(ErrorKind::Cancelled), TURN_CANCEL_MESSAGE),
216        }
217    }
218
219    fn start_llm_call(
220        &mut self,
221        call: LlmCallStart<'_>,
222        instrumentation: &OtelInstrumentation,
223        tools: &[ToolDefinition],
224        system_prompt: Option<&str>,
225    ) {
226        let model_name = call.model.unwrap_or(call.display_name).to_string();
227
228        // Metrics backends key series by attribute values, so the metric set
229        // stays a small fixed vocabulary — no content, no per-call details.
230        let mut metric_attributes = vec![
231            KeyValue::new(semconv::GEN_AI_OPERATION_NAME, "chat"),
232            KeyValue::new(semconv::GEN_AI_REQUEST_MODEL, model_name.clone()),
233        ];
234
235        if let Some(provider) = call.provider {
236            metric_attributes.push(KeyValue::new(semconv::GEN_AI_PROVIDER_NAME, genai_provider_name(provider)));
237        }
238
239        if call.purpose == LlmCallPurpose::Compaction {
240            metric_attributes.push(KeyValue::new(semconv::LLM_PURPOSE, "compaction"));
241        }
242
243        let mut attributes = metric_attributes.clone();
244        attributes.push(KeyValue::new(semconv::GEN_AI_REQUEST_STREAM, true));
245        attributes.push(KeyValue::new(semconv::LLM_ATTEMPT, i64::from(call.attempt)));
246        if let Some(pricing) = call.pricing {
247            attributes.extend(pricing_attributes(pricing));
248        }
249        // Only chat calls carry the turn's input and tool definitions; a
250        // compaction call's actual input is the internal summarization prompt.
251        if call.purpose == LlmCallPurpose::Chat {
252            if let Some(input) = self.input.get() {
253                attributes.extend(hashed_content(
254                    semconv::GEN_AI_INPUT_MESSAGES,
255                    semconv::AETHER_INPUT_MESSAGES_SHA256,
256                    input_messages_json(input),
257                    input,
258                ));
259            }
260            if instrumentation.content.tool_definitions && !tools.is_empty() {
261                attributes.push(KeyValue::new(semconv::GEN_AI_TOOL_DEFINITIONS, tool_definitions_json(tools)));
262            }
263            if instrumentation.content.system_instructions
264                && let Some(prompt) = system_prompt
265            {
266                attributes.extend(hashed_content(
267                    semconv::GEN_AI_SYSTEM_INSTRUCTIONS,
268                    semconv::AETHER_SYSTEM_INSTRUCTIONS_SHA256,
269                    system_instructions_json(prompt),
270                    prompt,
271                ));
272            }
273        }
274
275        let name = if model_name.is_empty() { "chat".to_string() } else { format!("chat {model_name}") };
276        let builder = SpanBuilder::from_name(name).with_kind(SpanKind::Client).with_attributes(attributes);
277        let context = instrumentation.start_span(builder, Some(self.span.context()));
278        let state = LlmCallState::new(
279            context,
280            instrumentation.metrics.clone(),
281            call.purpose,
282            instrumentation.content.output_messages,
283            metric_attributes,
284        );
285        *self.llm_call_slot(call.purpose) = Some(state);
286    }
287
288    fn on_tool_call(&mut self, request: &ToolCallRequest, instrumentation: &OtelInstrumentation) {
289        if let Some(chat) = &mut self.chat_call {
290            chat.record_tool_call_start(request);
291        }
292        let mut arguments = ContentBuffer::new(instrumentation.content.tool_calls);
293        arguments.set(&request.arguments);
294        self.streamed_arguments.insert(request.id.clone(), arguments);
295    }
296
297    fn on_tool_call_update(&mut self, tool_call_id: &str, chunk: &str) {
298        if let Some(chat) = &mut self.chat_call {
299            chat.record_tool_call_update(tool_call_id, chunk);
300        }
301        if let Some(arguments) = self.streamed_arguments.get_mut(tool_call_id) {
302            arguments.push(chunk);
303        }
304    }
305
306    fn on_tool_execution_started(&mut self, tool_id: &str, tool_name: &str, instrumentation: &OtelInstrumentation) {
307        let mut attributes = vec![
308            KeyValue::new(semconv::GEN_AI_OPERATION_NAME, "execute_tool"),
309            KeyValue::new(semconv::GEN_AI_TOOL_NAME, tool_name.to_string()),
310            KeyValue::new(semconv::GEN_AI_TOOL_CALL_ID, tool_id.to_string()),
311            KeyValue::new(semconv::MCP_METHOD_NAME, "tools/call"),
312            // The server sees the tool without our namespacing prefix, so this
313            // is what joins this span to the server's span for the same call.
314            KeyValue::new(semconv::MCP_TOOL_NAME, mcp_tool_name(tool_name).to_string()),
315        ];
316        let arguments = self.streamed_arguments.remove(tool_id);
317
318        if let Some(text) = arguments.as_ref().and_then(ContentBuffer::get) {
319            attributes.push(KeyValue::new(semconv::GEN_AI_TOOL_CALL_ARGUMENTS, text.to_string()));
320        }
321
322        let builder = SpanBuilder::from_name(format!("execute_tool {tool_name}"))
323            .with_kind(SpanKind::Client)
324            .with_attributes(attributes);
325        let context = instrumentation.start_span(builder, Some(self.span.context()));
326        self.executing_tools.insert(tool_id.to_string(), SpanGuard::new(context, TOOL_CANCEL_MESSAGE));
327    }
328
329    fn on_tool_result(&mut self, result: &ToolCallResult, instrumentation: &OtelInstrumentation) {
330        self.streamed_arguments.remove(&result.id);
331        let Some(mut span) = self.executing_tools.remove(&result.id) else { return };
332
333        if instrumentation.content.tool_calls {
334            span.set_attribute(KeyValue::new(semconv::GEN_AI_TOOL_CALL_RESULT, result.result.clone()));
335        }
336
337        span.end_ok();
338    }
339
340    fn on_tool_error(&mut self, error: &ToolCallError) {
341        self.streamed_arguments.remove(&error.id);
342        if let Some(mut span) = self.executing_tools.remove(&error.id) {
343            span.end_error(Some(ErrorKind::ToolError), error.error.clone());
344        }
345    }
346
347    fn llm_call_slot(&mut self, purpose: LlmCallPurpose) -> &mut Option<LlmCallState> {
348        match purpose {
349            LlmCallPurpose::Chat => &mut self.chat_call,
350            LlmCallPurpose::Compaction => &mut self.compaction_call,
351        }
352    }
353}
354
355/// Maps the agent's canonical provider name to the `GenAI` semantic-convention
356/// provider name, passing providers the catalog doesn't know through as-is.
357fn genai_provider_name(provider: &str) -> String {
358    provider.parse::<Provider>().map_or_else(|_| provider.to_string(), |p| p.genai_provider_name().to_string())
359}
360
361/// Borrowed view of [`TurnEvent::LlmCallStarted`]; named fields keep the two
362/// optional strings from being swapped at the call site.
363#[derive(Clone, Copy)]
364struct LlmCallStart<'a> {
365    purpose: LlmCallPurpose,
366    provider: Option<&'a str>,
367    model: Option<&'a str>,
368    display_name: &'a str,
369    pricing: Option<ModelPricing>,
370    attempt: u32,
371}
372
373fn hashed_content(content_key: &'static str, hash_key: &'static str, content: String, text: &str) -> [KeyValue; 2] {
374    [KeyValue::new(content_key, content), KeyValue::new(hash_key, sha256_hex(text))]
375}
376
377fn pricing_attributes(pricing: ModelPricing) -> Vec<KeyValue> {
378    let mut attributes = vec![
379        KeyValue::new(semconv::AI_INPUT_TOKEN_PRICE, pricing.input_per_million / TOKENS_PER_MILLION),
380        KeyValue::new(semconv::AI_OUTPUT_TOKEN_PRICE, pricing.output_per_million / TOKENS_PER_MILLION),
381    ];
382
383    if let Some(price) = pricing.cache_read_per_million {
384        attributes.push(KeyValue::new(semconv::AI_CACHE_READ_TOKEN_PRICE, price / TOKENS_PER_MILLION));
385    }
386
387    if let Some(price) = pricing.cache_write_per_million {
388        attributes.push(KeyValue::new(semconv::AI_CACHE_WRITE_TOKEN_PRICE, price / TOKENS_PER_MILLION));
389    }
390
391    attributes
392}
393
394const TOKENS_PER_MILLION: f64 = 1_000_000.0;