1use 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#[doc(hidden)]
21pub use tracing as __tracing;
22
23pub use tracing::field::Empty;
25
26#[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 { $($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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
89pub enum GenAiOperation {
90 Chat,
92 #[deprecated(note = "use `Chat`; streaming is recorded by `SpanBuilder::streaming`")]
94 ChatStreaming,
95 GenerateContent,
97 #[deprecated(note = "use `Chat`; streaming is recorded by `SpanBuilder::streaming`")]
99 Interactions,
100 #[deprecated(note = "use `Chat`; streaming is recorded by `SpanBuilder::streaming`")]
102 InteractionsStreaming,
103 Embeddings,
105 Rerank,
107 Transcription,
109 ImageGeneration,
111 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 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
148pub const PROVIDER_REQUEST_ID_FIELD: &str = "rig.provider_request_id";
154
155pub const COMPLETION_PARENT_MARKER_FIELD: &str = "rig.completion_parent";
160
161pub 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#[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 (
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
247pub use crate::completion_parent_span;
251
252#[derive(Debug, Clone, Copy, PartialEq, Eq)]
254enum CompletionParentVerdict {
255 Adopt,
257 RejectMissingFields,
261 NotAParent,
264}
265
266fn 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
276fn classify_completion_parent(metadata: &tracing::Metadata<'_>) -> CompletionParentVerdict {
278 let fields = metadata.fields();
279 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
295static NEAR_MISS_WARNED: LazyLock<Mutex<HashSet<Identifier>>> =
298 LazyLock::new(|| Mutex::new(HashSet::new()));
299
300#[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
316fn 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 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
369pub 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 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 pub fn streaming(mut self, streaming: bool) -> Self {
399 self.streaming = Some(streaming);
400 self
401 }
402
403 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 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 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(¤t);
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
579fn 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 "unknown"
686 };
687 vec![TelemetryOutputMessage {
688 role: "assistant",
689 parts: assistant_parts(content),
690 finish_reason,
691 }]
692}
693
694pub 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
708pub 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
724pub 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
741pub trait SpanCombinator {
743 fn record_token_usage(&self, usage: &Usage);
745
746 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 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;