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