Skip to main content

rig_core/telemetry/
mod.rs

1//! GenAI tracing spans, completion-parent adoption, and opt-in content recording.
2//!
3//! ```
4//! use rig_core::telemetry::{GenAiOperation, SpanBuilder};
5//!
6//! let span = SpanBuilder::new("provider", "model", GenAiOperation::Chat).build();
7//! ```
8use crate::completion::{AssistantContent, Message, Usage};
9use crate::message::{
10    DocumentSourceKind, Image, MimeType, Opaque, ToolResult, ToolResultContent, UserContent,
11};
12use base64::Engine;
13use serde::Serialize;
14use std::collections::HashSet;
15use std::sync::{LazyLock, Mutex};
16use tracing::callsite::Identifier;
17
18/// Macro implementation dependency; public because exported macro expansions
19/// must be able to resolve it from downstream crates.
20#[doc(hidden)]
21pub use tracing as __tracing;
22
23/// Declares a span field without a value, including in [`completion_parent_span!`].
24pub use tracing::field::Empty;
25
26/// Declares caller-supplied header fields and canonical completion fields.
27/// The header must end with a comma; nonempty extras must begin with one.
28#[doc(hidden)]
29#[macro_export]
30macro_rules! __rig_canonical_completion_span {
31    (
32        target: $target:literal,
33        $(parent: $parent:expr,)?
34        name: $name:literal,
35        // Both blocks are spliced verbatim into `info_span!`: the header block
36        // must end with a trailing comma, the extras block must begin with one.
37        // Violating either surfaces as an `info_span!` parse error at the call
38        // site, not here.
39        { $($header:tt)* }
40        { $($extra:tt)* }
41    ) => {
42        $crate::telemetry::__tracing::info_span!(
43            target: $target,
44            $(parent: $parent,)?
45            $name,
46            $($header)*
47            gen_ai.request.stream = $crate::telemetry::__tracing::field::Empty,
48            gen_ai.response.id = $crate::telemetry::__tracing::field::Empty,
49            gen_ai.response.model = $crate::telemetry::__tracing::field::Empty,
50            rig.provider_request_id = $crate::telemetry::__tracing::field::Empty,
51            gen_ai.usage.input_tokens = $crate::telemetry::__tracing::field::Empty,
52            gen_ai.usage.output_tokens = $crate::telemetry::__tracing::field::Empty,
53            gen_ai.usage.cache_read.input_tokens = $crate::telemetry::__tracing::field::Empty,
54            gen_ai.usage.cache_creation.input_tokens = $crate::telemetry::__tracing::field::Empty,
55            gen_ai.usage.tool_use_prompt_tokens = $crate::telemetry::__tracing::field::Empty,
56            gen_ai.usage.reasoning_tokens = $crate::telemetry::__tracing::field::Empty,
57            gen_ai.input.messages = $crate::telemetry::__tracing::field::Empty,
58            gen_ai.output.messages = $crate::telemetry::__tracing::field::Empty
59            $($extra)*
60        )
61    };
62}
63
64macro_rules! new_completion_span {
65    ($name:literal, $provider:expr, $request_model:expr, $operation:expr, $system:expr) => {
66        $crate::__rig_canonical_completion_span!(
67            target: "rig::completions",
68            name: $name,
69            {
70                gen_ai.operation.name = $operation,
71                gen_ai.provider.name = $provider,
72                gen_ai.request.model = $request_model,
73                gen_ai.system_instructions = $system,
74            }
75            {}
76        )
77    };
78}
79
80/// A GenAI operation and its canonical span name. Completion operations carry
81/// message content and may adopt a completion-parent span; the others record
82/// only usage and identity on a fresh span.
83///
84/// Completion operations report the well-known `gen_ai.operation.name` values
85/// of the OpenTelemetry GenAI semantic conventions. Streaming is a delivery
86/// mode, not an operation: it is recorded as `gen_ai.request.stream` (see
87/// [`SpanBuilder::streaming`]).
88#[derive(Debug, Clone, Copy, PartialEq, Eq)]
89pub enum GenAiOperation {
90    /// A chat completion (`chat`), unary or streamed.
91    Chat,
92    /// A streaming chat completion.
93    #[deprecated(note = "use `Chat`; streaming is recorded by `SpanBuilder::streaming`")]
94    ChatStreaming,
95    /// A Gemini generate-content request (`generate_content`), unary or streamed.
96    GenerateContent,
97    /// A Gemini Interactions API request.
98    #[deprecated(note = "use `Chat`; streaming is recorded by `SpanBuilder::streaming`")]
99    Interactions,
100    /// A streaming Gemini Interactions API request.
101    #[deprecated(note = "use `Chat`; streaming is recorded by `SpanBuilder::streaming`")]
102    InteractionsStreaming,
103    /// A text (or image) embedding request.
104    Embeddings,
105    /// A reranking request.
106    Rerank,
107    /// An audio transcription request.
108    Transcription,
109    /// An image generation request.
110    ImageGeneration,
111    /// An audio generation (text-to-speech) request.
112    AudioGeneration,
113}
114
115#[allow(deprecated)]
116impl GenAiOperation {
117    fn as_str(self) -> &'static str {
118        match self {
119            Self::Chat | Self::ChatStreaming | Self::Interactions | Self::InteractionsStreaming => {
120                "chat"
121            }
122            Self::GenerateContent => "generate_content",
123            Self::Embeddings => "embeddings",
124            Self::Rerank => "rerank",
125            Self::Transcription => "transcription",
126            Self::ImageGeneration => "image_generation",
127            Self::AudioGeneration => "audio_generation",
128        }
129    }
130
131    /// The delivery mode a deprecated streaming variant implies.
132    fn implied_streaming(self) -> Option<bool> {
133        matches!(self, Self::ChatStreaming | Self::InteractionsStreaming).then_some(true)
134    }
135
136    pub(crate) fn is_completion(self) -> bool {
137        matches!(
138            self,
139            Self::Chat
140                | Self::ChatStreaming
141                | Self::GenerateContent
142                | Self::Interactions
143                | Self::InteractionsStreaming
144        )
145    }
146}
147
148/// Span field holding the provider's transport request id (the request-id
149/// response header), recorded on success and on provider errors. GenAI
150/// semantic conventions define no attribute for it. Rig's own completion
151/// spans declare it; it is not required of an adopted parent, which simply
152/// does not record it when undeclared.
153pub const PROVIDER_REQUEST_ID_FIELD: &str = "rig.provider_request_id";
154
155/// Marker field for completion-parent adoption, independent of tracing target.
156/// Its value is ignored. Adoption requires every field in
157/// [`COMPLETION_PARENT_REQUIRED_FIELDS`]; otherwise the builder creates a child
158/// span and warns once per incomplete callsite.
159pub const COMPLETION_PARENT_MARKER_FIELD: &str = "rig.completion_parent";
160
161/// Fields required alongside [`COMPLETION_PARENT_MARKER_FIELD`] for adoption.
162/// Each must be statically declared with [`Empty`] or a value.
163/// [`completion_parent_span!`] declares the complete set.
164pub const COMPLETION_PARENT_REQUIRED_FIELDS: &[&str] = &[
165    "gen_ai.operation.name",
166    "gen_ai.provider.name",
167    "gen_ai.request.model",
168    "gen_ai.system_instructions",
169    "gen_ai.response.id",
170    "gen_ai.response.model",
171    "gen_ai.usage.input_tokens",
172    "gen_ai.usage.output_tokens",
173    "gen_ai.usage.cache_read.input_tokens",
174    "gen_ai.usage.cache_creation.input_tokens",
175    "gen_ai.usage.tool_use_prompt_tokens",
176    "gen_ai.usage.reasoning_tokens",
177    "gen_ai.input.messages",
178    "gen_ai.output.messages",
179];
180
181/// Declares a completion-parent span with the marker and all required fields.
182/// Defaults to the current span as parent. An optional `parent: <expr>` between
183/// `target` and `name` accepts any parent supported by [`tracing::info_span!`],
184/// including `None`.
185///
186/// Records the supplied operation and system instructions. Provider and model
187/// fields are populated on adoption. Extra fields must not duplicate the marker
188/// or required fields: recording a duplicate name updates only its first field.
189/// Use [`Empty`] or `Option::<&str>::None` for unset values without a direct
190/// `tracing` dependency.
191///
192/// ```
193/// use rig_core::telemetry::completion_parent_span;
194///
195/// let span = completion_parent_span!(
196///     target: "my_runtime",
197///     name: "chat",
198///     operation: "chat",
199///     system_instructions: Option::<&str>::None,
200///     gen_ai.agent.name = "assistant",
201/// );
202/// ```
203#[macro_export]
204macro_rules! completion_parent_span {
205    (
206        target: $target:literal,
207        parent: $parent:expr,
208        name: $name:literal,
209        operation: $operation:expr,
210        system_instructions: $system:expr
211        $(, $($extra:tt)*)?
212    ) => {
213        $crate::__rig_canonical_completion_span!(
214            target: $target,
215            parent: $parent,
216            name: $name,
217            {
218                rig.completion_parent = true,
219                gen_ai.operation.name = $operation,
220                gen_ai.system_instructions = $system,
221                gen_ai.provider.name = $crate::telemetry::__tracing::field::Empty,
222                gen_ai.request.model = $crate::telemetry::__tracing::field::Empty,
223            }
224            { $(, $($extra)*)? }
225        )
226    };
227    // Default arm: delegates to the explicit-parent arm so the two cannot
228    // drift in the fields they declare.
229    (
230        target: $target:literal,
231        name: $name:literal,
232        operation: $operation:expr,
233        system_instructions: $system:expr
234        $(, $($extra:tt)*)?
235    ) => {
236        $crate::completion_parent_span!(
237            target: $target,
238            parent: $crate::telemetry::__tracing::Span::current(),
239            name: $name,
240            operation: $operation,
241            system_instructions: $system
242            $(, $($extra)*)?
243        )
244    };
245}
246
247// `#[macro_export]` places the macro at the crate root; re-export it here so
248// it is also reachable at its documented home alongside the contract
249// constants it implements.
250pub use crate::completion_parent_span;
251
252/// Classification of the current span for completion-parent adoption.
253#[derive(Debug, Clone, Copy, PartialEq, Eq)]
254enum CompletionParentVerdict {
255    /// The marker and all required fields are present.
256    Adopt,
257    /// The marker, but the span omits at least one field in
258    /// [`COMPLETION_PARENT_REQUIRED_FIELDS`]. The missing names are computed
259    /// only if a warning is emitted.
260    RejectMissingFields,
261    /// No marker at all: an ordinary ambient span that becomes the parent of a
262    /// fresh `rig::completions` child. Never warns.
263    NotAParent,
264}
265
266/// Returns undeclared required fields for the warning diagnostic.
267fn missing_required_fields(metadata: &tracing::Metadata<'_>) -> Vec<&'static str> {
268    let fields = metadata.fields();
269    COMPLETION_PARENT_REQUIRED_FIELDS
270        .iter()
271        .copied()
272        .filter(|name| fields.field(name).is_none())
273        .collect()
274}
275
276/// Classifies span metadata without logging or changing global state.
277fn classify_completion_parent(metadata: &tracing::Metadata<'_>) -> CompletionParentVerdict {
278    let fields = metadata.fields();
279    // Exact match, never a prefix: a runtime field that merely starts with the
280    // marker name (`rig.completion_parent.id`, say) is not the marker and must
281    // not make its span a rejected parent.
282    if fields.field(COMPLETION_PARENT_MARKER_FIELD).is_none() {
283        return CompletionParentVerdict::NotAParent;
284    }
285    if COMPLETION_PARENT_REQUIRED_FIELDS
286        .iter()
287        .all(|name| fields.field(name).is_some())
288    {
289        CompletionParentVerdict::Adopt
290    } else {
291        CompletionParentVerdict::RejectMissingFields
292    }
293}
294
295/// Incomplete parent callsites already reported, bounded by the program's
296/// static callsites. Each callsite receives its own missing-field diagnostic.
297static NEAR_MISS_WARNED: LazyLock<Mutex<HashSet<Identifier>>> =
298    LazyLock::new(|| Mutex::new(HashSet::new()));
299
300/// Clear the per-callsite warn budget.
301///
302/// The budget is process-global, so without this the warning tests are coupled:
303/// whichever runs first consumes the budget for any callsite they share, and the
304/// other sees silence. `cargo nextest` hides that (one process per test) while
305/// `cargo test` exposes it, so the coupling would be green in CI and red
306/// locally — the worst orientation for a latent test bug. Resetting makes each
307/// test independent of callsite identity, ordering, and runner.
308#[cfg(test)]
309fn reset_near_miss_warnings() {
310    NEAR_MISS_WARNED
311        .lock()
312        .unwrap_or_else(|poisoned| poisoned.into_inner())
313        .clear();
314}
315
316/// Warns once per incomplete parent callsite, recovering a poisoned dedup lock.
317/// A subscriber panic must not prevent subsequent completions.
318fn warn_once_on_completion_parent_verdict(
319    verdict: CompletionParentVerdict,
320    metadata: &tracing::Metadata<'_>,
321) {
322    match verdict {
323        CompletionParentVerdict::Adopt | CompletionParentVerdict::NotAParent => {}
324        CompletionParentVerdict::RejectMissingFields => {
325            let first_sighting = {
326                let mut warned = NEAR_MISS_WARNED
327                    .lock()
328                    .unwrap_or_else(std::sync::PoisonError::into_inner);
329                warned.insert(metadata.callsite())
330            };
331            // Release the lock before invoking subscribers, which may re-enter
332            // the builder and deadlock if the guard is still held.
333            if !first_sighting {
334                return;
335            }
336            tracing::warn!(
337                marker = COMPLETION_PARENT_MARKER_FIELD,
338                missing_fields = ?missing_required_fields(metadata),
339                "completion-parent span declares the marker but not every required field \
340                 and is not adopted; provider telemetry lands on a fresh child span \
341                 instead — declare the span with \
342                 `rig_core::telemetry::completion_parent_span!`"
343            );
344        }
345    }
346}
347
348macro_rules! new_modality_span {
349    ($name:literal, $provider:expr, $request_model:expr, $operation:expr) => {
350        $crate::telemetry::__tracing::info_span!(
351            target: "rig::modalities",
352            $name,
353            gen_ai.operation.name = $operation,
354            gen_ai.provider.name = $provider,
355            gen_ai.request.model = $request_model,
356            gen_ai.response.id = $crate::telemetry::__tracing::field::Empty,
357            gen_ai.response.model = $crate::telemetry::__tracing::field::Empty,
358            rig.provider_request_id = $crate::telemetry::__tracing::field::Empty,
359            gen_ai.usage.input_tokens = $crate::telemetry::__tracing::field::Empty,
360            gen_ai.usage.output_tokens = $crate::telemetry::__tracing::field::Empty,
361            gen_ai.usage.cache_read.input_tokens = $crate::telemetry::__tracing::field::Empty,
362            gen_ai.usage.cache_creation.input_tokens = $crate::telemetry::__tracing::field::Empty,
363            gen_ai.usage.tool_use_prompt_tokens = $crate::telemetry::__tracing::field::Empty,
364            gen_ai.usage.reasoning_tokens = $crate::telemetry::__tracing::field::Empty,
365        )
366    };
367}
368
369/// Builder for a canonical GenAI span.
370///
371/// A completion operation reuses the current span when it declares
372/// [`COMPLETION_PARENT_MARKER_FIELD`] and every field in
373/// [`COMPLETION_PARENT_REQUIRED_FIELDS`], and otherwise opens a
374/// `rig::completions` child of the current span. Every other operation opens a
375/// fresh `rig::modalities` span.
376pub struct SpanBuilder<'a> {
377    provider: &'a str,
378    request_model: &'a str,
379    operation: GenAiOperation,
380    system_instructions: Option<String>,
381    streaming: Option<bool>,
382}
383
384impl<'a> SpanBuilder<'a> {
385    /// Create a span builder for a provider request.
386    pub fn new(provider: &'a str, request_model: &'a str, operation: GenAiOperation) -> Self {
387        Self {
388            provider,
389            request_model,
390            operation,
391            system_instructions: None,
392            streaming: operation.implied_streaming(),
393        }
394    }
395
396    /// Record whether the request streams its response, as
397    /// `gen_ai.request.stream`. Only completion spans record it.
398    pub fn streaming(mut self, streaming: bool) -> Self {
399        self.streaming = Some(streaming);
400        self
401    }
402
403    /// Set the system instructions sent with the request when sensitive content
404    /// telemetry has been explicitly enabled. Only completion spans record them.
405    pub fn system_instructions(
406        mut self,
407        system_instructions: Option<&'a str>,
408        record_content: bool,
409    ) -> Self {
410        self.system_instructions = system_instructions_json(system_instructions, record_content);
411        self
412    }
413
414    /// Build the operation's canonical span, or enrich Rig's current
415    /// completion-parent span for a completion operation.
416    pub fn build(self) -> tracing::Span {
417        if self.operation.is_completion()
418            && let Some(parent) = self.adopt_completion_parent()
419        {
420            return parent;
421        }
422
423        let (provider, model) = (self.provider, self.request_model);
424        let op = self.operation.as_str();
425        let sys = self.system_instructions.as_deref();
426        #[allow(deprecated)]
427        let span = match self.operation {
428            GenAiOperation::Chat
429            | GenAiOperation::ChatStreaming
430            | GenAiOperation::Interactions
431            | GenAiOperation::InteractionsStreaming => {
432                new_completion_span!("chat", provider, model, op, sys)
433            }
434            GenAiOperation::GenerateContent => {
435                new_completion_span!("generate_content", provider, model, op, sys)
436            }
437            GenAiOperation::Embeddings => new_modality_span!("embeddings", provider, model, op),
438            GenAiOperation::Rerank => new_modality_span!("rerank", provider, model, op),
439            GenAiOperation::Transcription => {
440                new_modality_span!("transcription", provider, model, op)
441            }
442            GenAiOperation::ImageGeneration => {
443                new_modality_span!("image_generation", provider, model, op)
444            }
445            GenAiOperation::AudioGeneration => {
446                new_modality_span!("audio_generation", provider, model, op)
447            }
448        };
449        if self.operation.is_completion() {
450            self.record_streaming(&span);
451        }
452        span
453    }
454
455    fn record_streaming(&self, span: &tracing::Span) {
456        if let Some(streaming) = self.streaming {
457            span.record("gen_ai.request.stream", streaming);
458        }
459    }
460
461    /// The current span, enriched with this request, when it is a conforming
462    /// completion parent. Warns once per near-miss callsite.
463    fn adopt_completion_parent(&self) -> Option<tracing::Span> {
464        let current = tracing::Span::current();
465        let metadata = current.metadata()?;
466        let verdict = classify_completion_parent(metadata);
467        warn_once_on_completion_parent_verdict(verdict, metadata);
468        if verdict != CompletionParentVerdict::Adopt {
469            return None;
470        }
471        current.record("gen_ai.operation.name", self.operation.as_str());
472        current.record("gen_ai.provider.name", self.provider);
473        current.record("gen_ai.request.model", self.request_model);
474        self.record_streaming(&current);
475        if let Some(system_instructions) = self.system_instructions.as_deref() {
476            current.record("gen_ai.system_instructions", system_instructions);
477        }
478        Some(current)
479    }
480}
481
482#[derive(Serialize)]
483struct TelemetryChatMessage {
484    role: &'static str,
485    parts: Vec<TelemetryPart>,
486}
487
488#[derive(Serialize)]
489struct TelemetryOutputMessage {
490    role: &'static str,
491    parts: Vec<TelemetryPart>,
492    finish_reason: &'static str,
493}
494
495#[derive(Serialize)]
496#[serde(tag = "type", rename_all = "snake_case")]
497enum TelemetryPart {
498    Text {
499        content: String,
500    },
501    ToolCall {
502        #[serde(skip_serializing_if = "Option::is_none")]
503        id: Option<String>,
504        name: String,
505        arguments: serde_json::Value,
506    },
507    ToolCallResponse {
508        #[serde(skip_serializing_if = "Option::is_none")]
509        id: Option<String>,
510        response: serde_json::Value,
511    },
512    Reasoning {
513        content: String,
514    },
515    Opaque {
516        kind: String,
517    },
518    Uri {
519        #[serde(skip_serializing_if = "Option::is_none")]
520        mime_type: Option<String>,
521        modality: &'static str,
522        uri: String,
523    },
524    File {
525        #[serde(skip_serializing_if = "Option::is_none")]
526        mime_type: Option<String>,
527        modality: &'static str,
528        file_id: String,
529    },
530    Blob {
531        #[serde(skip_serializing_if = "Option::is_none")]
532        mime_type: Option<String>,
533        modality: &'static str,
534        content: String,
535    },
536}
537
538fn media_part<T>(
539    data: &DocumentSourceKind,
540    media_type: Option<&T>,
541    modality: &'static str,
542) -> Option<TelemetryPart>
543where
544    T: MimeType,
545{
546    let mime_type = media_type.map(|media_type| media_type.to_mime_type().to_string());
547    match data {
548        DocumentSourceKind::Url(uri) => Some(TelemetryPart::Uri {
549            mime_type,
550            modality,
551            uri: uri.clone(),
552        }),
553        DocumentSourceKind::FileId(file_id) => Some(TelemetryPart::File {
554            mime_type,
555            modality,
556            file_id: file_id.clone(),
557        }),
558        DocumentSourceKind::Base64(content) => Some(TelemetryPart::Blob {
559            mime_type,
560            modality,
561            content: content.clone(),
562        }),
563        DocumentSourceKind::Raw(content) => Some(TelemetryPart::Blob {
564            mime_type,
565            modality,
566            content: base64::engine::general_purpose::STANDARD.encode(content),
567        }),
568        DocumentSourceKind::String(content) => Some(TelemetryPart::Text {
569            content: content.clone(),
570        }),
571        DocumentSourceKind::Unknown => None,
572    }
573}
574
575fn image_part(image: &Image) -> Option<TelemetryPart> {
576    media_part(&image.data, image.media_type.as_ref(), "image")
577}
578
579/// A provider-only item, named by its `type` and nothing else: its payload
580/// is the provider's, not content.
581fn opaque_part(opaque: &Opaque) -> TelemetryPart {
582    TelemetryPart::Opaque {
583        kind: opaque.kind().unwrap_or("unknown").to_owned(),
584    }
585}
586
587fn tool_result_response(result: &ToolResult) -> serde_json::Value {
588    let mut content = result
589        .content
590        .iter()
591        .filter_map(|content| match content {
592            ToolResultContent::Text(text) => Some(serde_json::Value::String(text.text.clone())),
593            ToolResultContent::Json { value } => Some(value.clone()),
594            ToolResultContent::Image(image) => {
595                image_part(image).and_then(|part| serde_json::to_value(part).ok())
596            }
597        })
598        .collect::<Vec<_>>();
599
600    if content.len() == 1 {
601        content.pop().unwrap_or(serde_json::Value::Null)
602    } else {
603        serde_json::Value::Array(content)
604    }
605}
606
607fn user_parts(content: &[UserContent]) -> Vec<TelemetryPart> {
608    content
609        .iter()
610        .filter_map(|content| match content {
611            UserContent::Text(text) => Some(TelemetryPart::Text {
612                content: text.text.clone(),
613            }),
614            UserContent::ToolResult(result) => Some(TelemetryPart::ToolCallResponse {
615                id: Some(result.call.to_string()),
616                response: tool_result_response(result),
617            }),
618            UserContent::Image(image) => image_part(image),
619            UserContent::Audio(audio) => {
620                media_part(&audio.data, audio.media_type.as_ref(), "audio")
621            }
622            UserContent::Video(video) => {
623                media_part(&video.data, video.media_type.as_ref(), "video")
624            }
625            UserContent::Document(document) => {
626                media_part(&document.data, document.media_type.as_ref(), "document")
627            }
628        })
629        .collect()
630}
631
632fn assistant_parts(content: &[AssistantContent]) -> Vec<TelemetryPart> {
633    content
634        .iter()
635        .flat_map(|content| match content {
636            AssistantContent::Text(text) => vec![TelemetryPart::Text {
637                content: text.text.clone(),
638            }],
639            AssistantContent::ToolCall(tool_call) => vec![TelemetryPart::ToolCall {
640                id: Some(tool_call.id.to_string()),
641                name: tool_call.function.name.clone().into(),
642                arguments: tool_call.function.arguments_value(),
643            }],
644            AssistantContent::Reasoning(reasoning) => vec![TelemetryPart::Reasoning {
645                content: reasoning.text.clone(),
646            }],
647            AssistantContent::Image(image) => image_part(image).into_iter().collect(),
648            AssistantContent::Opaque(opaque) => vec![opaque_part(opaque)],
649        })
650        .collect()
651}
652
653fn input_messages(messages: &[Message]) -> Vec<TelemetryChatMessage> {
654    messages
655        .iter()
656        .map(|message| match message {
657            Message::System { content } => TelemetryChatMessage {
658                role: "system",
659                parts: vec![TelemetryPart::Text {
660                    content: content.clone(),
661                }],
662            },
663            Message::User { content } => TelemetryChatMessage {
664                role: "user",
665                parts: user_parts(content),
666            },
667            Message::Assistant(turn) => TelemetryChatMessage {
668                role: "assistant",
669                parts: assistant_parts(&turn.content),
670            },
671        })
672        .collect()
673}
674
675fn output_messages(content: &[AssistantContent]) -> Vec<TelemetryOutputMessage> {
676    let finish_reason = if content
677        .iter()
678        .any(|content| matches!(content, AssistantContent::ToolCall(_)))
679    {
680        "tool_call"
681    } else {
682        // Rig's normalized assistant content does not retain provider finish
683        // reasons such as length or content filtering. Avoid claiming a clean
684        // stop when the actual reason is unavailable.
685        "unknown"
686    };
687    vec![TelemetryOutputMessage {
688        role: "assistant",
689        parts: assistant_parts(content),
690        finish_reason,
691    }]
692}
693
694/// Serializes system instructions using the normalized GenAI telemetry shape.
695pub fn system_instructions_json(instructions: Option<&str>, enabled: bool) -> Option<String> {
696    if !enabled {
697        return None;
698    }
699
700    instructions.and_then(|instructions| {
701        serde_json::to_string(&vec![TelemetryPart::Text {
702            content: instructions.to_string(),
703        }])
704        .ok()
705    })
706}
707
708/// Records serialized model input messages on `gen_ai.input.messages` when
709/// content telemetry is explicitly enabled.
710///
711/// Message content can contain prompts, retrieved context, tool results, and
712/// other sensitive or high-cardinality data. Keep this disabled unless the
713/// caller has explicitly opted in for debugging/observability.
714pub fn record_model_input(span: &tracing::Span, messages: &[Message], enabled: bool) {
715    if !enabled || span.is_disabled() {
716        return;
717    }
718
719    if let Ok(messages) = serde_json::to_string(&input_messages(messages)) {
720        span.record("gen_ai.input.messages", messages);
721    }
722}
723
724/// Records serialized model output messages on `gen_ai.output.messages` when
725/// content telemetry is explicitly enabled.
726///
727/// Message content can contain model responses, tool calls, and other sensitive
728/// or high-cardinality data. Keep this disabled unless the caller has explicitly
729/// opted in for debugging/observability.
730pub fn record_model_output(span: &tracing::Span, content: &[AssistantContent], enabled: bool) {
731    if !enabled || span.is_disabled() {
732        return;
733    }
734
735    let messages = output_messages(content);
736    if let Ok(messages) = serde_json::to_string(&messages) {
737        span.record("gen_ai.output.messages", messages);
738    }
739}
740
741/// Records GenAI usage and response metadata on tracing spans.
742pub trait SpanCombinator {
743    /// Record Rig-normalized token usage fields on the span.
744    fn record_token_usage(&self, usage: &Usage);
745
746    /// Record a response's ID, model, and token usage on the span.
747    fn record_response(&self, response_id: Option<&str>, model: Option<&str>, usage: &Usage);
748}
749
750impl SpanCombinator for tracing::Span {
751    fn record_token_usage(&self, usage: &Usage) {
752        if self.is_disabled() {
753            return;
754        }
755
756        // A counter the provider did not report leaves its span field unset;
757        // a reported zero is recorded as zero.
758        let fields = [
759            ("gen_ai.usage.input_tokens", usage.input_tokens),
760            ("gen_ai.usage.output_tokens", usage.output_tokens),
761            (
762                "gen_ai.usage.cache_read.input_tokens",
763                usage.cached_input_tokens,
764            ),
765            (
766                "gen_ai.usage.cache_creation.input_tokens",
767                usage.cache_creation_input_tokens,
768            ),
769            (
770                "gen_ai.usage.tool_use_prompt_tokens",
771                usage.tool_use_prompt_tokens,
772            ),
773            ("gen_ai.usage.reasoning_tokens", usage.reasoning_tokens),
774        ];
775        for (field, value) in fields {
776            if let Some(value) = value {
777                self.record(field, value);
778            }
779        }
780    }
781
782    fn record_response(&self, response_id: Option<&str>, model: Option<&str>, usage: &Usage) {
783        if self.is_disabled() {
784            return;
785        }
786        if let Some(id) = response_id {
787            self.record("gen_ai.response.id", id);
788        }
789        if let Some(model) = model {
790            self.record("gen_ai.response.model", model);
791        }
792        self.record_token_usage(usage);
793    }
794}
795
796#[cfg(test)]
797mod equivalence_tests;
798#[cfg(test)]
799mod tests;