Skip to main content

agent_top_core/
otlp.rs

1//! Spans an agent sends over OTLP, reduced to what agent-top keeps (RFC-107
2//! D3, DEC-026).
3//!
4//! The decoders in the binary turn a request into `ReceivedSpan`s, passing
5//! every attribute key through `keep`. Keys not on the allow-list are dropped
6//! there, so prompt, completion, system-instruction and tool argument or
7//! result attributes never reach the store, and a content attribute added
8//! upstream later is dropped by default.
9//!
10//! Attribute names follow the OpenTelemetry GenAI semantic conventions,
11//! verified against `semantic-conventions-genai` at schema
12//! `gen-ai-dev/1.42.0-dev` (commit cb10b70, 2026-10-05). In that version
13//! `gen_ai.usage.input_tokens` includes the cache read and cache write
14//! tokens, so the fresh input is what is left after subtracting them.
15
16use crate::model::{SpanKind, TokenUsage};
17use std::collections::BTreeMap;
18use std::time::{Duration, SystemTime, UNIX_EPOCH};
19
20/// Every attribute key a received span or resource may keep.
21const ALLOWED: &[&str] = &[
22    "service.name",
23    "gen_ai.operation.name",
24    "gen_ai.provider.name",
25    "gen_ai.request.model",
26    "gen_ai.response.model",
27    "gen_ai.conversation.id",
28    "gen_ai.agent.name",
29    "gen_ai.tool.name",
30    "gen_ai.usage.input_tokens",
31    "gen_ai.usage.output_tokens",
32    "gen_ai.usage.cache_read.input_tokens",
33    "gen_ai.usage.cache_write.input_tokens",
34    // Names older instrumentations still send.
35    "gen_ai.usage.prompt_tokens",
36    "gen_ai.usage.completion_tokens",
37    "gen_ai.system",
38    // What `agent-top trace --format otlp` writes.
39    "agent_top.harness",
40    "agent_top.session_id",
41    "agent_top.model",
42    "agent_top.kind",
43    "agent_top.sidechain",
44    "agent_top.mcp.server",
45    "agent_top.open",
46];
47
48/// Whether an attribute with this key is kept.
49pub fn keep(key: &str) -> bool {
50    ALLOWED.contains(&key)
51}
52
53/// An attribute value. Arrays, maps and bytes are never on the allow-list.
54#[derive(Debug, Clone, PartialEq)]
55pub enum Attr {
56    Str(String),
57    Int(i64),
58    Float(f64),
59    Bool(bool),
60}
61
62impl Attr {
63    pub fn as_str(&self) -> Option<&str> {
64        match self {
65            Attr::Str(s) => Some(s),
66            _ => None,
67        }
68    }
69
70    /// A count: an integer, or a string or float holding one.
71    pub fn as_u64(&self) -> Option<u64> {
72        match self {
73            Attr::Int(i) => u64::try_from(*i).ok(),
74            Attr::Float(f) if *f >= 0.0 && f.fract() == 0.0 => Some(*f as u64),
75            Attr::Str(s) => s.parse().ok(),
76            _ => None,
77        }
78    }
79
80    pub fn as_bool(&self) -> Option<bool> {
81        match self {
82            Attr::Bool(b) => Some(*b),
83            _ => None,
84        }
85    }
86}
87
88/// One received span, with only allow-listed attributes. Resource attributes
89/// are merged in under the span's own, which win.
90#[derive(Debug, Clone, PartialEq)]
91pub struct ReceivedSpan {
92    /// Lower-case hex.
93    pub trace_id: String,
94    pub span_id: String,
95    pub parent_span_id: Option<String>,
96    pub name: String,
97    pub start_unix_nano: u64,
98    pub end_unix_nano: u64,
99    pub error: bool,
100    pub attrs: BTreeMap<String, Attr>,
101}
102
103/// Where a span's tokens count: every inference span, and an agent span only
104/// in a session that sent no inference span with usage, so an agent span
105/// that repeats its children's totals is not counted twice.
106#[derive(Debug, Clone, Copy, PartialEq, Eq)]
107pub enum UsageRole {
108    Inference,
109    Agent,
110    None,
111}
112
113impl ReceivedSpan {
114    fn str(&self, key: &str) -> Option<&str> {
115        self.attrs.get(key).and_then(Attr::as_str)
116    }
117
118    fn count(&self, key: &str) -> Option<u64> {
119        self.attrs.get(key).and_then(Attr::as_u64)
120    }
121
122    pub fn operation(&self) -> Option<&str> {
123        self.str("gen_ai.operation.name")
124    }
125
126    /// The span's kind in agent-top's model, from `gen_ai.operation.name`,
127    /// else the `agent_top.kind` the trace export writes. Other operations
128    /// (embeddings, retrieval, memory) have none.
129    pub fn kind(&self) -> Option<SpanKind> {
130        match self.operation() {
131            Some("chat" | "generate_content" | "text_completion") => Some(SpanKind::Inference),
132            Some("execute_tool") => Some(SpanKind::Tool),
133            Some("invoke_agent" | "invoke_workflow") => Some(SpanKind::Turn),
134            Some(_) => None,
135            None => match self.str("agent_top.kind") {
136                Some("inference") => Some(SpanKind::Inference),
137                Some("tool") => Some(SpanKind::Tool),
138                Some("turn") => Some(SpanKind::Turn),
139                _ => None,
140            },
141        }
142    }
143
144    pub fn usage_role(&self) -> UsageRole {
145        match self.operation() {
146            Some("chat" | "generate_content" | "text_completion" | "embeddings") => UsageRole::Inference,
147            Some("invoke_agent" | "invoke_workflow") => UsageRole::Agent,
148            _ => UsageRole::None,
149        }
150    }
151
152    /// The name a span is shown under: the tool for a tool call, else the
153    /// span's own name.
154    pub fn display_name(&self) -> String {
155        match (self.kind(), self.str("gen_ai.tool.name")) {
156            (Some(SpanKind::Tool), Some(t)) => t.to_string(),
157            _ => self.name.clone(),
158        }
159    }
160
161    /// The response model, else the requested one, else the one the trace
162    /// export names.
163    pub fn model(&self) -> Option<&str> {
164        self.str("gen_ai.response.model").or_else(|| self.str("gen_ai.request.model")).or_else(|| self.str("agent_top.model"))
165    }
166
167    pub fn provider(&self) -> Option<&str> {
168        self.str("gen_ai.provider.name").or_else(|| self.str("gen_ai.system"))
169    }
170
171    pub fn service(&self) -> &str {
172        self.str("service.name").unwrap_or("unknown_service")
173    }
174
175    /// `agent_top.harness`, else `otel`.
176    pub fn harness(&self) -> &str {
177        self.str("agent_top.harness").unwrap_or("otel")
178    }
179
180    /// The conversation this span belongs to: `gen_ai.conversation.id`, else
181    /// the trace export's session id, else the trace id.
182    pub fn conversation(&self) -> &str {
183        self.str("gen_ai.conversation.id").or_else(|| self.str("agent_top.session_id")).unwrap_or(&self.trace_id)
184    }
185
186    /// The store's session id for this span: `<service.name>/<conversation>`,
187    /// so a received session never collides with a transcript's id.
188    pub fn session_id(&self) -> String {
189        format!("{}/{}", self.service(), self.conversation())
190    }
191
192    pub fn sidechain(&self) -> bool {
193        self.attrs.get("agent_top.sidechain").and_then(Attr::as_bool).unwrap_or(false)
194    }
195
196    /// Tokens, split the way the price table charges them.
197    pub fn usage(&self) -> TokenUsage {
198        let cache_read = self.count("gen_ai.usage.cache_read.input_tokens").unwrap_or(0);
199        let cache_write = self.count("gen_ai.usage.cache_write.input_tokens").unwrap_or(0);
200        let input = self.count("gen_ai.usage.input_tokens").or_else(|| self.count("gen_ai.usage.prompt_tokens")).unwrap_or(0);
201        let output = self.count("gen_ai.usage.output_tokens").or_else(|| self.count("gen_ai.usage.completion_tokens")).unwrap_or(0);
202        TokenUsage {
203            input: input.saturating_sub(cache_read + cache_write),
204            // The semconv has one cache-write count with no TTL; the 5-minute
205            // rate is the default write every provider in the table uses.
206            cache_write_5m: cache_write,
207            cache_read,
208            output,
209            ..Default::default()
210        }
211    }
212
213    pub fn started_at(&self) -> SystemTime {
214        UNIX_EPOCH + Duration::from_nanos(self.start_unix_nano)
215    }
216
217    pub fn duration_ms(&self) -> u64 {
218        self.end_unix_nano.saturating_sub(self.start_unix_nano) / 1_000_000
219    }
220
221    /// The trace export marks a span that had not finished.
222    pub fn is_open(&self) -> bool {
223        self.attrs.get("agent_top.open").and_then(Attr::as_bool).unwrap_or(false)
224    }
225}
226
227#[cfg(test)]
228mod tests {
229    use super::*;
230
231    fn span(attrs: &[(&str, Attr)]) -> ReceivedSpan {
232        ReceivedSpan {
233            trace_id: "t".into(),
234            span_id: "s".into(),
235            parent_span_id: None,
236            name: "chat claude".into(),
237            start_unix_nano: 1_000_000_000,
238            end_unix_nano: 1_250_000_000,
239            error: false,
240            attrs: attrs.iter().map(|(k, v)| (k.to_string(), v.clone())).collect(),
241        }
242    }
243
244    #[test]
245    fn content_attributes_are_not_on_the_allow_list() {
246        for k in [
247            "gen_ai.input.messages",
248            "gen_ai.output.messages",
249            "gen_ai.system_instructions",
250            "gen_ai.tool.call.arguments",
251            "gen_ai.tool.call.result",
252            "gen_ai.prompt.0.content",
253        ] {
254            assert!(!keep(k), "{k}");
255        }
256        assert!(keep("gen_ai.usage.input_tokens"));
257    }
258
259    #[test]
260    fn input_tokens_include_the_cache_and_are_split_out() {
261        let s = span(&[
262            ("gen_ai.operation.name", Attr::Str("chat".into())),
263            ("gen_ai.usage.input_tokens", Attr::Int(165)),
264            ("gen_ai.usage.cache_read.input_tokens", Attr::Int(100)),
265            ("gen_ai.usage.cache_write.input_tokens", Attr::Int(25)),
266            ("gen_ai.usage.output_tokens", Attr::Str("12".into())),
267        ]);
268        let u = s.usage();
269        assert_eq!((u.input, u.cache_read, u.cache_write_5m, u.output), (40, 100, 25, 12));
270        assert_eq!(u.total(), 177);
271        assert_eq!(s.kind(), Some(SpanKind::Inference));
272        assert_eq!(s.usage_role(), UsageRole::Inference);
273        assert_eq!(s.duration_ms(), 250);
274    }
275
276    #[test]
277    fn older_names_still_count() {
278        let s = span(&[("gen_ai.usage.prompt_tokens", Attr::Int(10)), ("gen_ai.usage.completion_tokens", Attr::Int(3))]);
279        assert_eq!((s.usage().input, s.usage().output), (10, 3));
280    }
281
282    #[test]
283    fn kinds_come_from_the_operation_or_the_trace_export() {
284        let tool = span(&[("gen_ai.operation.name", Attr::Str("execute_tool".into())), ("gen_ai.tool.name", Attr::Str("search".into()))]);
285        assert_eq!((tool.kind(), tool.display_name().as_str()), (Some(SpanKind::Tool), "search"));
286        assert_eq!(span(&[("gen_ai.operation.name", Attr::Str("invoke_agent".into()))]).kind(), Some(SpanKind::Turn));
287        assert_eq!(span(&[("gen_ai.operation.name", Attr::Str("embeddings".into()))]).kind(), None);
288        assert_eq!(span(&[("agent_top.kind", Attr::Str("tool".into()))]).kind(), Some(SpanKind::Tool));
289        assert_eq!(span(&[]).kind(), None);
290    }
291
292    #[test]
293    fn a_session_is_keyed_by_service_and_conversation() {
294        let s = span(&[("service.name", Attr::Str("triage-bot".into())), ("gen_ai.conversation.id", Attr::Str("c1".into()))]);
295        assert_eq!(s.session_id(), "triage-bot/c1");
296        assert_eq!(span(&[]).session_id(), "unknown_service/t");
297        assert_eq!(s.harness(), "otel");
298    }
299}