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
19pub struct OtelObserver {
22 turn: Option<TurnState>,
23 tool_definitions: Vec<ToolDefinition>,
24 system_prompt: Option<String>,
25 otel: OtelInstrumentation,
26}
27
28#[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 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 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
120struct TurnState {
125 span: SpanGuard,
126 input: ContentBuffer,
127 output: ContentBuffer,
128 chat_call: Option<LlmCallState>,
129 compaction_call: Option<LlmCallState>,
130 streamed_arguments: HashMap<String, ContentBuffer>,
132 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 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 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 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 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
355fn 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#[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;