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 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 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 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 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 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
337fn 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#[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;