1use 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#[doc(hidden)]
22pub use tracing as __tracing;
23
24pub use tracing::field::Empty;
26
27#[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 { $($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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
84pub enum GenAiOperation {
85 Chat,
87 ChatStreaming,
89 GenerateContent,
91 Interactions,
93 InteractionsStreaming,
95 Embeddings,
97 Rerank,
99 Transcription,
101 ImageGeneration,
103 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
135pub const PROVIDER_REQUEST_ID_FIELD: &str = "rig.provider_request_id";
141
142pub const COMPLETION_PARENT_MARKER_FIELD: &str = "rig.completion_parent";
147
148pub 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#[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 (
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
234pub use crate::completion_parent_span;
238
239#[derive(Debug, Clone, Copy, PartialEq, Eq)]
241enum CompletionParentVerdict {
242 Adopt,
244 RejectMissingFields,
248 NotAParent,
251}
252
253fn 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
263fn classify_completion_parent(metadata: &tracing::Metadata<'_>) -> CompletionParentVerdict {
265 let fields = metadata.fields();
266 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
282static NEAR_MISS_WARNED: LazyLock<Mutex<HashSet<Identifier>>> =
285 LazyLock::new(|| Mutex::new(HashSet::new()));
286
287#[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
303fn 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 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
356pub 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 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 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 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 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 "unknown"
659 };
660 vec![TelemetryOutputMessage {
661 role: "assistant",
662 parts: assistant_parts(content),
663 finish_reason,
664 }]
665}
666
667pub 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
681pub 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
697pub 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
714pub trait SpanCombinator {
716 fn record_token_usage(&self, usage: &Usage);
718
719 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 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;