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
19pub struct OtelObserver {
22 turn: Option<TurnState>,
23 tool_definitions: Vec<ToolDefinition>,
24 otel: OtelInstrumentation,
25}
26
27#[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 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 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
113struct TurnState {
118 span: SpanGuard,
119 input: ContentBuffer,
120 output: ContentBuffer,
121 chat_call: Option<LlmCallState>,
122 compaction_call: Option<LlmCallState>,
123 streamed_arguments: HashMap<String, ContentBuffer>,
125 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 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 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 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 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
333fn 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#[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;