Skip to main content

openrouter/
webhooks.rs

1//! Parser for OpenRouter Broadcast webhook payloads (OTLP JSON traces).
2//!
3//! Use [`parse_broadcast_payload`] for the raw OTLP envelope,
4//! [`extract_broadcast_traces`] to flatten it into [`BroadcastTrace`] rows,
5//! or [`parse_broadcast_traces`] to do both in one call. The convenience
6//! function is what most webhook handlers want.
7//!
8//! Wiring this into an HTTP handler is intentionally left to the caller —
9//! the parser is framework-agnostic. A typical axum handler looks like:
10//!
11//! ```ignore
12//! async fn webhook(body: axum::body::Bytes) -> impl axum::response::IntoResponse {
13//!     match openrouter::webhooks::parse_broadcast_traces(&body) {
14//!         Ok(traces) => { /* process */ axum::http::StatusCode::OK }
15//!         Err(_)     => axum::http::StatusCode::BAD_REQUEST,
16//!     }
17//! }
18//! ```
19//!
20//! Shapes mirror the Go SDK (`broadcast.go`, `broadcast_models.go`).
21
22use 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/// Top-level OTLP JSON trace payload sent by Broadcast.
31#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
32pub struct OtlpExportTraceRequest {
33    /// Resource-grouped span batches.
34    #[serde(rename = "resourceSpans", default)]
35    pub resource_spans: Vec<OtlpResourceSpan>,
36}
37
38/// Spans grouped by their originating resource.
39#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
40pub struct OtlpResourceSpan {
41    /// The resource describing the trace origin.
42    #[serde(default)]
43    pub resource: OtlpResource,
44    /// Spans grouped by instrumentation scope.
45    #[serde(rename = "scopeSpans", default)]
46    pub scope_spans: Vec<OtlpScopeSpan>,
47}
48
49/// The entity producing telemetry.
50#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
51pub struct OtlpResource {
52    /// Resource-level attributes.
53    #[serde(default)]
54    pub attributes: Vec<OtlpAttribute>,
55}
56
57/// Spans grouped by instrumentation scope.
58#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
59pub struct OtlpScopeSpan {
60    /// Identifies the instrumentation library, when set.
61    #[serde(default)]
62    pub scope: Option<OtlpScope>,
63    /// Spans in this scope.
64    #[serde(default)]
65    pub spans: Vec<OtlpSpan>,
66}
67
68/// Identifies the instrumentation library.
69#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
70pub struct OtlpScope {
71    /// Library name.
72    #[serde(default)]
73    pub name: String,
74    /// Library version.
75    #[serde(default)]
76    pub version: String,
77}
78
79/// A single span within a trace.
80#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
81pub struct OtlpSpan {
82    /// Hex trace identifier.
83    #[serde(rename = "traceId", default)]
84    pub trace_id: String,
85    /// Hex span identifier.
86    #[serde(rename = "spanId", default)]
87    pub span_id: String,
88    /// Parent span identifier (empty for the root).
89    #[serde(rename = "parentSpanId", default)]
90    pub parent_span_id: String,
91    /// Span name.
92    #[serde(default)]
93    pub name: String,
94    /// OTLP span kind code (1 = internal, 2 = server, ...).
95    #[serde(default)]
96    pub kind: i32,
97    /// Start time as nanoseconds since the Unix epoch (string).
98    #[serde(rename = "startTimeUnixNano", default)]
99    pub start_time_unix_nano: String,
100    /// End time as nanoseconds since the Unix epoch (string).
101    #[serde(rename = "endTimeUnixNano", default)]
102    pub end_time_unix_nano: String,
103    /// Attributes attached to the span.
104    #[serde(default)]
105    pub attributes: Vec<OtlpAttribute>,
106    /// Span status.
107    #[serde(default)]
108    pub status: Option<OtlpStatus>,
109    /// Span events.
110    #[serde(default)]
111    pub events: Vec<OtlpEvent>,
112}
113
114/// Key-value pair attached to a span or resource.
115#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
116pub struct OtlpAttribute {
117    /// Attribute name.
118    #[serde(default)]
119    pub key: String,
120    /// Attribute value.
121    #[serde(default)]
122    pub value: OtlpAnyValue,
123}
124
125/// A polymorphic OTLP value. The OTLP spec encodes int64 values as strings,
126/// but some emitters send them as JSON numbers — the SDK tolerates both
127/// via the internal `deserialize_flex_int_opt` helper.
128#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
129pub struct OtlpAnyValue {
130    /// String-typed value.
131    #[serde(
132        rename = "stringValue",
133        default,
134        skip_serializing_if = "Option::is_none"
135    )]
136    pub string_value: Option<String>,
137    /// Int64-typed value, encoded as a string per OTLP spec.
138    #[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    /// Double-typed value.
146    #[serde(
147        rename = "doubleValue",
148        default,
149        skip_serializing_if = "Option::is_none"
150    )]
151    pub double_value: Option<f64>,
152    /// Boolean-typed value.
153    #[serde(rename = "boolValue", default, skip_serializing_if = "Option::is_none")]
154    pub bool_value: Option<bool>,
155    /// Array-typed value.
156    #[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    /// Render the value as a string, regardless of its underlying type.
181    /// Returns the empty string for [`Self::array_value`] or an
182    /// all-`None` value.
183    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            // Match Go's `%g` formatting reasonably well.
192            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/// Wraps a slice of OTLP values.
202#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
203pub struct OtlpArrayValue {
204    /// Element values.
205    #[serde(default)]
206    pub values: Vec<OtlpAnyValue>,
207}
208
209/// The status of a span.
210#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
211pub struct OtlpStatus {
212    /// OTLP status code (0 = unset, 1 = ok, 2 = error).
213    #[serde(default)]
214    pub code: i32,
215    /// Human-readable error message.
216    #[serde(default)]
217    pub message: String,
218}
219
220/// A timed event within a span.
221#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
222pub struct OtlpEvent {
223    /// Event name.
224    #[serde(default)]
225    pub name: String,
226    /// Event time as nanoseconds since the Unix epoch (string).
227    #[serde(rename = "timeUnixNano", default)]
228    pub time_unix_nano: String,
229    /// Event-specific attributes.
230    #[serde(default)]
231    pub attributes: Vec<OtlpAttribute>,
232}
233
234/// User-friendly representation of a single span extracted from an OTLP
235/// trace payload sent by OpenRouter Broadcast. Field names mirror the Go
236/// SDK; deprecated aliases ([`Self::prompt_tokens`], etc.) are kept so
237/// callers porting from the Go SDK don't have to rename everything.
238#[derive(Clone, Debug, Default, PartialEq)]
239pub struct BroadcastTrace {
240    /// Hex trace identifier.
241    pub trace_id: String,
242    /// Hex span identifier.
243    pub span_id: String,
244    /// Parent span identifier (empty for the root).
245    pub parent_span_id: String,
246    /// Span name.
247    pub span_name: String,
248
249    /// Start time as nanoseconds since the Unix epoch (0 when missing).
250    pub start_time_unix_nano: i64,
251    /// End time as nanoseconds since the Unix epoch (0 when missing).
252    pub end_time_unix_nano: i64,
253    /// `end - start` when both timestamps are present.
254    pub duration: Duration,
255
256    /// Deprecated: prompt tokens (use [`Self::input_tokens`]).
257    pub prompt_tokens: i64,
258    /// Deprecated: completion tokens (use [`Self::output_tokens`]).
259    pub completion_tokens: i64,
260    /// Total tokens (computed as input + output when absent upstream).
261    pub total_tokens: i64,
262    /// Total cost in USD (alias of [`Self::total_cost`]).
263    pub cost: f64,
264    /// Resolved model id (prefers `response.model` over `request.model`).
265    pub model: String,
266
267    /// Canonical input-token count.
268    pub input_tokens: i64,
269    /// Canonical output-token count.
270    pub output_tokens: i64,
271
272    /// Total USD cost.
273    pub total_cost: f64,
274    /// Input-side USD cost.
275    pub input_cost: f64,
276    /// Output-side USD cost.
277    pub output_cost: f64,
278
279    /// Tokens served from cache.
280    pub cached_tokens: i64,
281    /// Audio-input tokens.
282    pub audio_input_tokens: i64,
283    /// Video-input tokens.
284    pub video_input_tokens: i64,
285    /// Image-output tokens.
286    pub image_output_tokens: i64,
287    /// Reasoning tokens.
288    pub reasoning_tokens: i64,
289
290    /// `gen_ai.operation.name` value.
291    pub operation_name: String,
292    /// `gen_ai.system` value.
293    pub system: String,
294    /// `gen_ai.provider.name` value.
295    pub provider_name: String,
296    /// `gen_ai.response.model` value.
297    pub response_model: String,
298    /// `gen_ai.response.finish_reason` value.
299    pub finish_reason: String,
300    /// `gen_ai.response.finish_reasons` value.
301    pub finish_reasons: String,
302    /// `gen_ai.request.model` value.
303    pub request_model: String,
304
305    /// `openrouter.provider_slug` value.
306    pub provider_slug: String,
307    /// `openrouter.provider_name` value.
308    pub openrouter_provider_name: String,
309    /// `openrouter.api_key_name` value.
310    pub api_key_name: String,
311    /// `openrouter.entity_id` value.
312    pub entity_id: String,
313    /// `openrouter.user_id` value.
314    pub openrouter_user_id: String,
315    /// `openrouter.finish_reason` value.
316    pub openrouter_finish_reason: String,
317    /// Per-input-token USD unit price.
318    pub input_unit_price: f64,
319    /// Per-output-token USD unit price.
320    pub output_unit_price: f64,
321    /// `openrouter.source` value (originating client/integration).
322    pub source: String,
323
324    /// Verbatim `gen_ai.prompt` payload.
325    pub prompt: String,
326    /// Verbatim `gen_ai.completion` payload.
327    pub completion: String,
328
329    /// `span.type` value.
330    pub span_type: String,
331    /// `span.level` value.
332    pub span_level: String,
333    /// `span.input` value.
334    pub span_input: String,
335    /// `span.output` value.
336    pub span_output: String,
337
338    /// `trace.name` value.
339    pub trace_name: String,
340    /// `trace.input` value.
341    pub trace_input: String,
342    /// `trace.output` value.
343    pub trace_output: String,
344    /// `trace.tags` value.
345    pub trace_tags: String,
346
347    /// `user.id` value.
348    pub user_id: String,
349    /// `session.id` value.
350    pub session_id: String,
351
352    /// Values from `trace.metadata.*` attributes (prefix stripped).
353    pub metadata: BTreeMap<String, String>,
354    /// Values from `span.metadata.*` attributes (prefix stripped).
355    pub span_metadata: BTreeMap<String, String>,
356    /// Attributes from the OTLP resource.
357    pub resource_attributes: BTreeMap<String, String>,
358    /// All other span attributes not mapped to a named field.
359    pub raw_attributes: BTreeMap<String, String>,
360}
361
362/// Parse raw JSON bytes into the OTLP trace envelope.
363pub fn parse_broadcast_payload(data: &[u8]) -> Result<OtlpExportTraceRequest> {
364    serde_json::from_slice(data).map_err(Error::Decode)
365}
366
367/// Flatten an OTLP envelope into a vector of [`BroadcastTrace`] rows.
368/// Missing attributes produce zero values; extraction is best-effort.
369pub 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
382/// Convenience: parse + flatten in one call.
383pub 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        // Model fields
420        "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        // Token usage (new canonical keys)
431        "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        // Token usage (old keys, backward compat)
440        "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        // Cost fields
456        "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        // Token detail
470        "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        // GenAI semantic convention
476        "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-specific
482        "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        // Content
492        "gen_ai.prompt" => t.prompt = val,
493        "gen_ai.completion" => t.completion = val,
494        // Span-level
495        "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-level
500        "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        // Identity
505        "user.id" => t.user_id = val,
506        "session.id" => t.session_id = val,
507        // Prefixed metadata / unknown attributes
508        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); // backward-compat alias populated
575        assert_eq!(t.output_tokens, 30);
576        assert_eq!(t.total_tokens, 150); // computed when absent
577        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); // back-filled
614        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}