1use std::collections::BTreeMap;
23use std::time::Duration;
24
25use serde::{Deserialize, Deserializer, Serialize};
26use serde_json::Value;
27
28use crate::error::{Error, Result};
29
30#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
32pub struct OtlpExportTraceRequest {
33 #[serde(rename = "resourceSpans", default)]
35 pub resource_spans: Vec<OtlpResourceSpan>,
36}
37
38#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
40pub struct OtlpResourceSpan {
41 #[serde(default)]
43 pub resource: OtlpResource,
44 #[serde(rename = "scopeSpans", default)]
46 pub scope_spans: Vec<OtlpScopeSpan>,
47}
48
49#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
51pub struct OtlpResource {
52 #[serde(default)]
54 pub attributes: Vec<OtlpAttribute>,
55}
56
57#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
59pub struct OtlpScopeSpan {
60 #[serde(default)]
62 pub scope: Option<OtlpScope>,
63 #[serde(default)]
65 pub spans: Vec<OtlpSpan>,
66}
67
68#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
70pub struct OtlpScope {
71 #[serde(default)]
73 pub name: String,
74 #[serde(default)]
76 pub version: String,
77}
78
79#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
81pub struct OtlpSpan {
82 #[serde(rename = "traceId", default)]
84 pub trace_id: String,
85 #[serde(rename = "spanId", default)]
87 pub span_id: String,
88 #[serde(rename = "parentSpanId", default)]
90 pub parent_span_id: String,
91 #[serde(default)]
93 pub name: String,
94 #[serde(default)]
96 pub kind: i32,
97 #[serde(rename = "startTimeUnixNano", default)]
99 pub start_time_unix_nano: String,
100 #[serde(rename = "endTimeUnixNano", default)]
102 pub end_time_unix_nano: String,
103 #[serde(default)]
105 pub attributes: Vec<OtlpAttribute>,
106 #[serde(default)]
108 pub status: Option<OtlpStatus>,
109 #[serde(default)]
111 pub events: Vec<OtlpEvent>,
112}
113
114#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
116pub struct OtlpAttribute {
117 #[serde(default)]
119 pub key: String,
120 #[serde(default)]
122 pub value: OtlpAnyValue,
123}
124
125#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
129pub struct OtlpAnyValue {
130 #[serde(
132 rename = "stringValue",
133 default,
134 skip_serializing_if = "Option::is_none"
135 )]
136 pub string_value: Option<String>,
137 #[serde(
139 rename = "intValue",
140 default,
141 skip_serializing_if = "Option::is_none",
142 deserialize_with = "deserialize_flex_int_opt"
143 )]
144 pub int_value: Option<String>,
145 #[serde(
147 rename = "doubleValue",
148 default,
149 skip_serializing_if = "Option::is_none"
150 )]
151 pub double_value: Option<f64>,
152 #[serde(rename = "boolValue", default, skip_serializing_if = "Option::is_none")]
154 pub bool_value: Option<bool>,
155 #[serde(
157 rename = "arrayValue",
158 default,
159 skip_serializing_if = "Option::is_none"
160 )]
161 pub array_value: Option<OtlpArrayValue>,
162}
163
164fn deserialize_flex_int_opt<'de, D>(de: D) -> std::result::Result<Option<String>, D::Error>
165where
166 D: Deserializer<'de>,
167{
168 let v = Option::<Value>::deserialize(de)?;
169 match v {
170 None | Some(Value::Null) => Ok(None),
171 Some(Value::String(s)) => Ok(Some(s)),
172 Some(Value::Number(n)) => Ok(Some(n.to_string())),
173 Some(other) => Err(serde::de::Error::custom(format!(
174 "intValue: expected string or number, got {other}"
175 ))),
176 }
177}
178
179impl OtlpAnyValue {
180 pub fn string_val(&self) -> String {
184 if let Some(s) = &self.string_value {
185 return s.clone();
186 }
187 if let Some(i) = &self.int_value {
188 return i.clone();
189 }
190 if let Some(d) = self.double_value {
191 return format!("{d}");
193 }
194 if let Some(b) = self.bool_value {
195 return b.to_string();
196 }
197 String::new()
198 }
199}
200
201#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
203pub struct OtlpArrayValue {
204 #[serde(default)]
206 pub values: Vec<OtlpAnyValue>,
207}
208
209#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
211pub struct OtlpStatus {
212 #[serde(default)]
214 pub code: i32,
215 #[serde(default)]
217 pub message: String,
218}
219
220#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
222pub struct OtlpEvent {
223 #[serde(default)]
225 pub name: String,
226 #[serde(rename = "timeUnixNano", default)]
228 pub time_unix_nano: String,
229 #[serde(default)]
231 pub attributes: Vec<OtlpAttribute>,
232}
233
234#[derive(Clone, Debug, Default, PartialEq)]
239pub struct BroadcastTrace {
240 pub trace_id: String,
242 pub span_id: String,
244 pub parent_span_id: String,
246 pub span_name: String,
248
249 pub start_time_unix_nano: i64,
251 pub end_time_unix_nano: i64,
253 pub duration: Duration,
255
256 pub prompt_tokens: i64,
258 pub completion_tokens: i64,
260 pub total_tokens: i64,
262 pub cost: f64,
264 pub model: String,
266
267 pub input_tokens: i64,
269 pub output_tokens: i64,
271
272 pub total_cost: f64,
274 pub input_cost: f64,
276 pub output_cost: f64,
278
279 pub cached_tokens: i64,
281 pub audio_input_tokens: i64,
283 pub video_input_tokens: i64,
285 pub image_output_tokens: i64,
287 pub reasoning_tokens: i64,
289
290 pub operation_name: String,
292 pub system: String,
294 pub provider_name: String,
296 pub response_model: String,
298 pub finish_reason: String,
300 pub finish_reasons: String,
302 pub request_model: String,
304
305 pub provider_slug: String,
307 pub openrouter_provider_name: String,
309 pub api_key_name: String,
311 pub entity_id: String,
313 pub openrouter_user_id: String,
315 pub openrouter_finish_reason: String,
317 pub input_unit_price: f64,
319 pub output_unit_price: f64,
321 pub source: String,
323
324 pub prompt: String,
326 pub completion: String,
328
329 pub span_type: String,
331 pub span_level: String,
333 pub span_input: String,
335 pub span_output: String,
337
338 pub trace_name: String,
340 pub trace_input: String,
342 pub trace_output: String,
344 pub trace_tags: String,
346
347 pub user_id: String,
349 pub session_id: String,
351
352 pub metadata: BTreeMap<String, String>,
354 pub span_metadata: BTreeMap<String, String>,
356 pub resource_attributes: BTreeMap<String, String>,
358 pub raw_attributes: BTreeMap<String, String>,
360}
361
362pub fn parse_broadcast_payload(data: &[u8]) -> Result<OtlpExportTraceRequest> {
364 serde_json::from_slice(data).map_err(Error::Decode)
365}
366
367pub fn extract_broadcast_traces(payload: &OtlpExportTraceRequest) -> Vec<BroadcastTrace> {
370 let mut out = Vec::new();
371 for rs in &payload.resource_spans {
372 let res_attrs = extract_attribute_map(&rs.resource.attributes);
373 for ss in &rs.scope_spans {
374 for span in &ss.spans {
375 out.push(build_trace(span, &res_attrs));
376 }
377 }
378 }
379 out
380}
381
382pub fn parse_broadcast_traces(data: &[u8]) -> Result<Vec<BroadcastTrace>> {
384 Ok(extract_broadcast_traces(&parse_broadcast_payload(data)?))
385}
386
387fn build_trace(span: &OtlpSpan, res_attrs: &BTreeMap<String, String>) -> BroadcastTrace {
388 let mut t = BroadcastTrace {
389 trace_id: span.trace_id.clone(),
390 span_id: span.span_id.clone(),
391 parent_span_id: span.parent_span_id.clone(),
392 span_name: span.name.clone(),
393 resource_attributes: res_attrs.clone(),
394 ..Default::default()
395 };
396 let start = span.start_time_unix_nano.parse::<i64>().unwrap_or(0);
397 let end = span.end_time_unix_nano.parse::<i64>().unwrap_or(0);
398 t.start_time_unix_nano = start;
399 t.end_time_unix_nano = end;
400 if start > 0 && end > start {
401 let delta = (end - start) as u64;
402 t.duration = Duration::from_nanos(delta);
403 }
404 for attr in &span.attributes {
405 let val = attr.value.string_val();
406 apply_attribute(&mut t, &attr.key, val);
407 }
408 if t.total_tokens == 0 && (t.input_tokens > 0 || t.output_tokens > 0) {
409 t.total_tokens = t.input_tokens + t.output_tokens;
410 }
411 t
412}
413
414fn apply_attribute(t: &mut BroadcastTrace, key: &str, val: String) {
415 let v = val.as_str();
416 let parse_i = |s: &str| s.parse::<i64>().unwrap_or(0);
417 let parse_f = |s: &str| s.parse::<f64>().unwrap_or(0.0);
418 match key {
419 "gen_ai.response.model" => {
421 t.response_model = val.clone();
422 t.model = val;
423 }
424 "gen_ai.request.model" => {
425 t.request_model = val.clone();
426 if t.model.is_empty() {
427 t.model = val;
428 }
429 }
430 "gen_ai.usage.input_tokens" => {
432 t.input_tokens = parse_i(v);
433 t.prompt_tokens = t.input_tokens;
434 }
435 "gen_ai.usage.output_tokens" => {
436 t.output_tokens = parse_i(v);
437 t.completion_tokens = t.output_tokens;
438 }
439 "gen_ai.usage.prompt_tokens" => {
441 let n = parse_i(v);
442 t.prompt_tokens = n;
443 if t.input_tokens == 0 {
444 t.input_tokens = n;
445 }
446 }
447 "gen_ai.usage.completion_tokens" => {
448 let n = parse_i(v);
449 t.completion_tokens = n;
450 if t.output_tokens == 0 {
451 t.output_tokens = n;
452 }
453 }
454 "gen_ai.usage.total_tokens" => t.total_tokens = parse_i(v),
455 "gen_ai.usage.total_cost" => {
457 t.total_cost = parse_f(v);
458 t.cost = t.total_cost;
459 }
460 "gen_ai.usage.cost" => {
461 let f = parse_f(v);
462 t.cost = f;
463 if t.total_cost == 0.0 {
464 t.total_cost = f;
465 }
466 }
467 "gen_ai.usage.input_cost" => t.input_cost = parse_f(v),
468 "gen_ai.usage.output_cost" => t.output_cost = parse_f(v),
469 "gen_ai.usage.input_tokens.cached" => t.cached_tokens = parse_i(v),
471 "gen_ai.usage.input_tokens.audio" => t.audio_input_tokens = parse_i(v),
472 "gen_ai.usage.input_tokens.video" => t.video_input_tokens = parse_i(v),
473 "gen_ai.usage.output_tokens.image" => t.image_output_tokens = parse_i(v),
474 "gen_ai.usage.output_tokens.reasoning" => t.reasoning_tokens = parse_i(v),
475 "gen_ai.operation.name" => t.operation_name = val,
477 "gen_ai.system" => t.system = val,
478 "gen_ai.provider.name" => t.provider_name = val,
479 "gen_ai.response.finish_reason" => t.finish_reason = val,
480 "gen_ai.response.finish_reasons" => t.finish_reasons = val,
481 "openrouter.provider_slug" => t.provider_slug = val,
483 "openrouter.provider_name" => t.openrouter_provider_name = val,
484 "openrouter.api_key_name" => t.api_key_name = val,
485 "openrouter.entity_id" => t.entity_id = val,
486 "openrouter.user_id" => t.openrouter_user_id = val,
487 "openrouter.finish_reason" => t.openrouter_finish_reason = val,
488 "openrouter.input_unit_price" => t.input_unit_price = parse_f(v),
489 "openrouter.output_unit_price" => t.output_unit_price = parse_f(v),
490 "openrouter.source" => t.source = val,
491 "gen_ai.prompt" => t.prompt = val,
493 "gen_ai.completion" => t.completion = val,
494 "span.type" => t.span_type = val,
496 "span.level" => t.span_level = val,
497 "span.input" => t.span_input = val,
498 "span.output" => t.span_output = val,
499 "trace.name" => t.trace_name = val,
501 "trace.input" => t.trace_input = val,
502 "trace.output" => t.trace_output = val,
503 "trace.tags" => t.trace_tags = val,
504 "user.id" => t.user_id = val,
506 "session.id" => t.session_id = val,
507 other => {
509 if let Some(rest) = other.strip_prefix("trace.metadata.") {
510 t.metadata.insert(rest.to_string(), val);
511 } else if let Some(rest) = other.strip_prefix("span.metadata.") {
512 t.span_metadata.insert(rest.to_string(), val);
513 } else {
514 t.raw_attributes.insert(other.to_string(), val);
515 }
516 }
517 }
518}
519
520fn extract_attribute_map(attrs: &[OtlpAttribute]) -> BTreeMap<String, String> {
521 attrs
522 .iter()
523 .map(|a| (a.key.clone(), a.value.string_val()))
524 .collect()
525}
526
527#[cfg(test)]
528mod tests {
529 use super::*;
530
531 #[test]
532 fn flex_int_accepts_string_or_number() {
533 let from_string: OtlpAnyValue = serde_json::from_str(r#"{"intValue":"42"}"#).unwrap();
534 assert_eq!(from_string.int_value.as_deref(), Some("42"));
535
536 let from_number: OtlpAnyValue = serde_json::from_str(r#"{"intValue":42}"#).unwrap();
537 assert_eq!(from_number.int_value.as_deref(), Some("42"));
538 }
539
540 #[test]
541 fn parse_minimal_payload_extracts_one_trace() {
542 let payload = br#"{
543 "resourceSpans": [{
544 "resource": {"attributes": [{"key":"service.name","value":{"stringValue":"or"}}]},
545 "scopeSpans": [{
546 "scope": {"name":"or-gateway","version":"1.0"},
547 "spans": [{
548 "traceId":"abc","spanId":"s1","name":"gen_ai.chat",
549 "kind":2,
550 "startTimeUnixNano":"1700000000000000000",
551 "endTimeUnixNano":"1700000000500000000",
552 "attributes": [
553 {"key":"gen_ai.response.model","value":{"stringValue":"openai/gpt-5"}},
554 {"key":"gen_ai.usage.input_tokens","value":{"intValue":"120"}},
555 {"key":"gen_ai.usage.output_tokens","value":{"intValue":"30"}},
556 {"key":"gen_ai.usage.total_cost","value":{"stringValue":"0.0042"}},
557 {"key":"openrouter.provider_slug","value":{"stringValue":"openai"}},
558 {"key":"trace.metadata.tenant","value":{"stringValue":"acme"}},
559 {"key":"span.metadata.region","value":{"stringValue":"us-west"}},
560 {"key":"some.custom.attr","value":{"stringValue":"x"}}
561 ]
562 }]
563 }]
564 }]
565 }"#;
566 let traces = parse_broadcast_traces(payload).unwrap();
567 assert_eq!(traces.len(), 1);
568 let t = &traces[0];
569 assert_eq!(t.trace_id, "abc");
570 assert_eq!(t.span_id, "s1");
571 assert_eq!(t.model, "openai/gpt-5");
572 assert_eq!(t.response_model, "openai/gpt-5");
573 assert_eq!(t.input_tokens, 120);
574 assert_eq!(t.prompt_tokens, 120); assert_eq!(t.output_tokens, 30);
576 assert_eq!(t.total_tokens, 150); assert!((t.total_cost - 0.0042).abs() < 1e-9);
578 assert!((t.cost - 0.0042).abs() < 1e-9);
579 assert_eq!(t.provider_slug, "openai");
580 assert_eq!(t.metadata.get("tenant").map(String::as_str), Some("acme"));
581 assert_eq!(
582 t.span_metadata.get("region").map(String::as_str),
583 Some("us-west")
584 );
585 assert_eq!(
586 t.raw_attributes.get("some.custom.attr").map(String::as_str),
587 Some("x")
588 );
589 assert_eq!(
590 t.resource_attributes
591 .get("service.name")
592 .map(String::as_str),
593 Some("or")
594 );
595 assert_eq!(t.duration, Duration::from_millis(500));
596 }
597
598 #[test]
599 fn old_token_keys_are_back_filled() {
600 let payload = br#"{
601 "resourceSpans":[{"resource":{"attributes":[]},"scopeSpans":[{"spans":[{
602 "traceId":"a","spanId":"b","name":"x",
603 "kind":1,"startTimeUnixNano":"0","endTimeUnixNano":"0",
604 "attributes":[
605 {"key":"gen_ai.usage.prompt_tokens","value":{"intValue":"10"}},
606 {"key":"gen_ai.usage.completion_tokens","value":{"intValue":"5"}}
607 ]
608 }]}]}]
609 }"#;
610 let traces = parse_broadcast_traces(payload).unwrap();
611 let t = &traces[0];
612 assert_eq!(t.prompt_tokens, 10);
613 assert_eq!(t.input_tokens, 10); assert_eq!(t.completion_tokens, 5);
615 assert_eq!(t.output_tokens, 5);
616 assert_eq!(t.total_tokens, 15);
617 }
618
619 #[test]
620 fn invalid_json_errors() {
621 let err = parse_broadcast_traces(b"not json").unwrap_err();
622 assert!(matches!(err, Error::Decode(_)));
623 }
624}