1use crate::value::VmDictExt;
19
20use crate::value::VmValue;
21
22pub mod source {
26 pub const OLLAMA_CHAT: &str = "ollama_chat";
28 pub const OLLAMA_GENERATE: &str = "ollama_generate";
30 pub const OPENAI_USAGE: &str = "openai_usage";
34 pub const LLAMACPP_TIMINGS: &str = "llamacpp_timings";
38 pub const ANTHROPIC_USAGE: &str = "anthropic_usage";
40 pub const GEMINI_USAGE: &str = "gemini_usage";
42 pub const GEMINI_INTERACTIONS_USAGE: &str = "gemini_interactions_usage";
46 pub const UNKNOWN: &str = "unknown";
50}
51
52pub(crate) fn elapsed_ms(started: std::time::Instant) -> u64 {
53 started.elapsed().as_millis().min(u128::from(u64::MAX)) as u64
54}
55
56#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
62pub struct ProviderTelemetry {
63 #[serde(default, skip_serializing_if = "String::is_empty")]
67 pub source: String,
68 #[serde(skip_serializing_if = "Option::is_none")]
71 pub serving_base_url: Option<String>,
72 #[serde(skip_serializing_if = "Option::is_none")]
74 pub server_total_ms: Option<u64>,
75 #[serde(skip_serializing_if = "Option::is_none")]
78 pub server_load_ms: Option<u64>,
79 #[serde(skip_serializing_if = "Option::is_none")]
83 pub server_prompt_eval_ms: Option<u64>,
84 #[serde(skip_serializing_if = "Option::is_none")]
86 pub server_generation_ms: Option<u64>,
87 #[serde(skip_serializing_if = "Option::is_none")]
92 pub server_prompt_tokens: Option<i64>,
93 #[serde(skip_serializing_if = "Option::is_none")]
95 pub server_output_tokens: Option<i64>,
96 #[serde(skip_serializing_if = "Option::is_none")]
100 pub client_wall_ms: Option<u64>,
101 #[serde(skip_serializing_if = "Option::is_none")]
104 pub runtime_context_length: Option<u64>,
105 #[serde(skip_serializing_if = "Option::is_none")]
108 pub runtime_loaded_model: Option<String>,
109 #[serde(skip_serializing_if = "Option::is_none")]
113 pub response_model: Option<String>,
114 #[serde(skip_serializing_if = "Option::is_none")]
116 pub runtime_memory_bytes: Option<u64>,
117 #[serde(skip_serializing_if = "Option::is_none")]
119 pub runtime_memory_vram_bytes: Option<u64>,
120 #[serde(skip_serializing_if = "Option::is_none")]
122 pub runtime_keep_alive_until: Option<String>,
123 #[serde(skip_serializing_if = "Option::is_none")]
125 pub request_id: Option<String>,
126 #[serde(skip_serializing_if = "Option::is_none")]
131 pub provider_metadata: Option<serde_json::Value>,
132}
133
134impl ProviderTelemetry {
135 pub fn new(source: &str) -> Self {
136 Self {
137 source: source.to_string(),
138 ..Self::default()
139 }
140 }
141
142 pub fn is_empty(&self) -> bool {
146 let Self {
147 source,
148 serving_base_url,
149 server_total_ms,
150 server_load_ms,
151 server_prompt_eval_ms,
152 server_generation_ms,
153 server_prompt_tokens,
154 server_output_tokens,
155 client_wall_ms,
156 runtime_context_length,
157 runtime_loaded_model,
158 response_model,
159 runtime_memory_bytes,
160 runtime_memory_vram_bytes,
161 runtime_keep_alive_until,
162 request_id,
163 provider_metadata,
164 } = self;
165 source.is_empty()
166 && serving_base_url.is_none()
167 && server_total_ms.is_none()
168 && server_load_ms.is_none()
169 && server_prompt_eval_ms.is_none()
170 && server_generation_ms.is_none()
171 && server_prompt_tokens.is_none()
172 && server_output_tokens.is_none()
173 && client_wall_ms.is_none()
174 && runtime_context_length.is_none()
175 && runtime_loaded_model.is_none()
176 && response_model.is_none()
177 && runtime_memory_bytes.is_none()
178 && runtime_memory_vram_bytes.is_none()
179 && runtime_keep_alive_until.is_none()
180 && request_id.is_none()
181 && provider_metadata.is_none()
182 }
183
184 pub fn ns_to_ms(ns: u64) -> u64 {
188 ns / 1_000_000
192 }
193
194 pub fn from_ollama_done(frame: &serde_json::Value, source: &str) -> Self {
198 let mut telemetry = Self::new(source);
199 telemetry.server_total_ms = ns_field(frame, "total_duration");
200 telemetry.server_load_ms = ns_field(frame, "load_duration");
201 telemetry.server_prompt_eval_ms = ns_field(frame, "prompt_eval_duration");
202 telemetry.server_generation_ms = ns_field(frame, "eval_duration");
203 telemetry.server_prompt_tokens = frame
204 .get("prompt_eval_count")
205 .and_then(serde_json::Value::as_i64);
206 telemetry.server_output_tokens =
207 frame.get("eval_count").and_then(serde_json::Value::as_i64);
208 if let Some(model) = frame.get("model").and_then(serde_json::Value::as_str) {
209 telemetry.runtime_loaded_model = Some(model.to_string());
210 }
211 telemetry
212 }
213
214 pub fn from_openai_usage(usage: &serde_json::Value, request_id: Option<&str>) -> Self {
218 let mut telemetry = Self::new(source::OPENAI_USAGE);
219 telemetry.server_prompt_tokens = usage
220 .get("prompt_tokens")
221 .or_else(|| usage.get("input_tokens"))
222 .and_then(serde_json::Value::as_i64);
223 telemetry.server_output_tokens = usage
224 .get("completion_tokens")
225 .or_else(|| usage.get("output_tokens"))
226 .and_then(serde_json::Value::as_i64);
227 if let Some(timings) = usage.get("timings").filter(|value| value.is_object()) {
228 telemetry.source = source::LLAMACPP_TIMINGS.to_string();
229 telemetry.server_prompt_eval_ms = ms_or_round(timings.get("prompt_ms"));
230 telemetry.server_generation_ms = ms_or_round(timings.get("predicted_ms"));
231 if let Some(prefill) = timings.get("prompt_n").and_then(serde_json::Value::as_i64) {
236 telemetry.server_prompt_tokens = Some(prefill);
237 }
238 if let Some(predicted) = timings
239 .get("predicted_n")
240 .and_then(serde_json::Value::as_i64)
241 {
242 telemetry.server_output_tokens = Some(predicted);
243 }
244 let total = telemetry
245 .server_prompt_eval_ms
246 .unwrap_or(0)
247 .saturating_add(telemetry.server_generation_ms.unwrap_or(0));
248 if total > 0 {
249 telemetry.server_total_ms = Some(total);
250 }
251 }
252 if let Some(request_id) = request_id.filter(|value| !value.is_empty()) {
253 telemetry.request_id = Some(request_id.to_string());
254 }
255 telemetry
256 }
257
258 pub fn from_anthropic_usage(usage: &serde_json::Value, request_id: Option<&str>) -> Self {
263 let mut telemetry = Self::new(source::ANTHROPIC_USAGE);
264 telemetry.server_prompt_tokens = usage
265 .get("input_tokens")
266 .and_then(serde_json::Value::as_i64);
267 telemetry.server_output_tokens = usage
268 .get("output_tokens")
269 .and_then(serde_json::Value::as_i64);
270 if let Some(request_id) = request_id.filter(|value| !value.is_empty()) {
271 telemetry.request_id = Some(request_id.to_string());
272 }
273 telemetry
274 }
275
276 pub fn from_gemini_usage(usage: &serde_json::Value, request_id: Option<&str>) -> Self {
280 let mut telemetry = Self::new(source::GEMINI_USAGE);
281 telemetry.server_prompt_tokens = usage
282 .get("promptTokenCount")
283 .and_then(serde_json::Value::as_i64);
284 telemetry.server_output_tokens = usage
285 .get("candidatesTokenCount")
286 .and_then(serde_json::Value::as_i64);
287 if let Some(request_id) = request_id.filter(|value| !value.is_empty()) {
288 telemetry.request_id = Some(request_id.to_string());
289 }
290 telemetry
291 }
292
293 pub fn from_gemini_interactions_usage(
298 usage: &serde_json::Value,
299 request_id: Option<&str>,
300 ) -> Self {
301 let mut telemetry = Self::new(source::GEMINI_INTERACTIONS_USAGE);
302 telemetry.server_prompt_tokens = usage
303 .get("total_input_tokens")
304 .and_then(serde_json::Value::as_i64);
305 telemetry.server_output_tokens = usage
306 .get("total_output_tokens")
307 .and_then(serde_json::Value::as_i64);
308 if let Some(request_id) = request_id.filter(|value| !value.is_empty()) {
309 telemetry.request_id = Some(request_id.to_string());
310 }
311 telemetry
312 }
313
314 pub fn capture_provider_metadata(&mut self, response: &serde_json::Value) {
317 if let Some(model) = response
318 .get("model")
319 .and_then(serde_json::Value::as_str)
320 .filter(|value| !value.is_empty())
321 {
322 self.response_model = Some(model.to_string());
323 }
324 if let Some(metadata) = response
325 .get("provider_metadata")
326 .filter(|value| !value.is_null())
327 .filter(|value| !value.as_object().is_some_and(serde_json::Map::is_empty))
328 {
329 merge_provider_metadata(&mut self.provider_metadata, metadata);
330 }
331 if let Some(receipt) = response
332 .get(crate::llm::managed_supply::MANAGED_SUPPLY_WIRE_KEY)
333 .filter(|value| !value.is_null())
334 {
335 let mut managed = serde_json::Map::new();
336 managed.insert(
337 crate::llm::managed_supply::MANAGED_SUPPLY_WIRE_KEY.to_string(),
338 receipt.clone(),
339 );
340 merge_provider_metadata(
341 &mut self.provider_metadata,
342 &serde_json::Value::Object(managed),
343 );
344 }
345 }
346
347 pub fn merge_ollama_ps(&mut self, ps: &OllamaPsModel) {
351 if self.runtime_loaded_model.is_none() {
352 self.runtime_loaded_model = ps.name.clone();
353 }
354 if self.runtime_context_length.is_none() {
355 self.runtime_context_length = ps.context_length;
356 }
357 if self.runtime_memory_bytes.is_none() {
358 self.runtime_memory_bytes = ps.size_bytes;
359 }
360 if self.runtime_memory_vram_bytes.is_none() {
361 self.runtime_memory_vram_bytes = ps.size_vram_bytes;
362 }
363 if self.runtime_keep_alive_until.is_none() {
364 self.runtime_keep_alive_until = ps.expires_at.clone();
365 }
366 }
367
368 pub fn as_vm_dict(&self) -> Option<VmValue> {
372 if self.is_empty() {
373 return None;
374 }
375 let mut dict: crate::value::DictMap = crate::value::DictMap::new();
376 if !self.source.is_empty() {
377 dict.put_str("source", self.source.as_str());
378 }
379 if let Some(ref serving_base_url) = self.serving_base_url {
380 dict.put_str("serving_base_url", serving_base_url.as_str());
381 }
382 insert_opt_u64(&mut dict, "server_total_ms", self.server_total_ms);
383 insert_opt_u64(&mut dict, "server_load_ms", self.server_load_ms);
384 insert_opt_u64(
385 &mut dict,
386 "server_prompt_eval_ms",
387 self.server_prompt_eval_ms,
388 );
389 insert_opt_u64(&mut dict, "server_generation_ms", self.server_generation_ms);
390 insert_opt_i64(&mut dict, "server_prompt_tokens", self.server_prompt_tokens);
391 insert_opt_i64(&mut dict, "server_output_tokens", self.server_output_tokens);
392 insert_opt_u64(&mut dict, "client_wall_ms", self.client_wall_ms);
393 insert_opt_u64(
394 &mut dict,
395 "runtime_context_length",
396 self.runtime_context_length,
397 );
398 if let Some(ref model) = self.runtime_loaded_model {
399 dict.put_str("runtime_loaded_model", model.as_str());
400 }
401 if let Some(ref model) = self.response_model {
402 dict.put_str("response_model", model.as_str());
403 }
404 insert_opt_u64(&mut dict, "runtime_memory_bytes", self.runtime_memory_bytes);
405 insert_opt_u64(
406 &mut dict,
407 "runtime_memory_vram_bytes",
408 self.runtime_memory_vram_bytes,
409 );
410 if let Some(ref expires) = self.runtime_keep_alive_until {
411 dict.put_str("runtime_keep_alive_until", expires.as_str());
412 }
413 if let Some(ref request_id) = self.request_id {
414 dict.put_str("request_id", request_id.as_str());
415 }
416 if let Some(ref provider_metadata) = self.provider_metadata {
417 dict.insert(
418 crate::value::intern_key("provider_metadata"),
419 crate::stdlib::json_to_vm_value(provider_metadata),
420 );
421 }
422 Some(VmValue::dict(dict))
423 }
424}
425
426fn merge_provider_metadata(target: &mut Option<serde_json::Value>, incoming: &serde_json::Value) {
427 let Some(incoming) = incoming.as_object() else {
428 return;
429 };
430 let target = target.get_or_insert_with(|| serde_json::Value::Object(Default::default()));
431 let Some(target) = target.as_object_mut() else {
432 return;
433 };
434 for (key, value) in incoming {
435 target.insert(key.clone(), value.clone());
436 }
437}
438
439#[derive(Clone, Debug, Default, PartialEq, Eq)]
442pub struct OllamaPsModel {
443 pub name: Option<String>,
444 pub size_bytes: Option<u64>,
445 pub size_vram_bytes: Option<u64>,
446 pub expires_at: Option<String>,
447 pub context_length: Option<u64>,
448}
449
450impl OllamaPsModel {
451 pub fn from_ps_entry(entry: &serde_json::Value) -> Option<Self> {
455 let name = entry
456 .get("name")
457 .and_then(serde_json::Value::as_str)
458 .or_else(|| entry.get("model").and_then(serde_json::Value::as_str))
459 .map(str::to_string);
460 let context_length = entry
461 .get("context_length")
462 .and_then(serde_json::Value::as_u64)
463 .or_else(|| {
464 entry
465 .get("details")
466 .and_then(|details| details.get("context_length"))
467 .and_then(serde_json::Value::as_u64)
468 });
469 Some(Self {
470 name,
471 size_bytes: entry.get("size").and_then(serde_json::Value::as_u64),
472 size_vram_bytes: entry.get("size_vram").and_then(serde_json::Value::as_u64),
473 expires_at: entry
474 .get("expires_at")
475 .and_then(serde_json::Value::as_str)
476 .map(str::to_string),
477 context_length,
478 })
479 }
480}
481
482fn ns_field(frame: &serde_json::Value, key: &str) -> Option<u64> {
483 frame
484 .get(key)
485 .and_then(serde_json::Value::as_u64)
486 .map(ProviderTelemetry::ns_to_ms)
487}
488
489fn ms_or_round(value: Option<&serde_json::Value>) -> Option<u64> {
490 let value = value?;
491 if let Some(n) = value.as_u64() {
492 return Some(n);
493 }
494 value.as_f64().map(|n| n.round().max(0.0) as u64)
495}
496
497fn insert_opt_u64(dict: &mut crate::value::DictMap, key: &str, value: Option<u64>) {
498 if let Some(value) = value {
499 dict.insert(crate::value::intern_key(key), VmValue::Int(value as i64));
500 }
501}
502
503fn insert_opt_i64(dict: &mut crate::value::DictMap, key: &str, value: Option<i64>) {
504 if let Some(value) = value {
505 dict.insert(crate::value::intern_key(key), VmValue::Int(value));
506 }
507}
508
509#[cfg(test)]
510mod tests {
511 use super::*;
512
513 #[test]
514 fn ollama_done_frame_extracts_full_breakdown() {
515 let frame = serde_json::json!({
516 "model": "devstral-small-2:24b",
517 "total_duration": 7_400_000_000u64,
518 "load_duration": 400_000_000u64,
519 "prompt_eval_duration": 1_200_000_000u64,
520 "eval_duration": 5_800_000_000u64,
521 "prompt_eval_count": 1024,
522 "eval_count": 64
523 });
524
525 let telemetry = ProviderTelemetry::from_ollama_done(&frame, source::OLLAMA_CHAT);
526
527 assert_eq!(telemetry.source, source::OLLAMA_CHAT);
528 assert_eq!(telemetry.server_total_ms, Some(7400));
529 assert_eq!(telemetry.server_load_ms, Some(400));
530 assert_eq!(telemetry.server_prompt_eval_ms, Some(1200));
531 assert_eq!(telemetry.server_generation_ms, Some(5800));
532 assert_eq!(telemetry.server_prompt_tokens, Some(1024));
533 assert_eq!(telemetry.server_output_tokens, Some(64));
534 assert_eq!(
535 telemetry.runtime_loaded_model.as_deref(),
536 Some("devstral-small-2:24b")
537 );
538 assert!(!telemetry.is_empty());
539 }
540
541 #[test]
542 fn ollama_done_frame_leaves_missing_fields_as_none() {
543 let frame = serde_json::json!({
544 "model": "devstral-small-2:24b",
545 });
547
548 let telemetry = ProviderTelemetry::from_ollama_done(&frame, source::OLLAMA_CHAT);
549
550 assert_eq!(telemetry.server_total_ms, None);
551 assert_eq!(telemetry.server_load_ms, None);
552 assert_eq!(telemetry.server_prompt_eval_ms, None);
553 assert_eq!(telemetry.server_generation_ms, None);
554 assert_eq!(telemetry.server_prompt_tokens, None);
555 assert_eq!(telemetry.server_output_tokens, None);
556 }
557
558 #[test]
559 fn openai_usage_partial_extracts_counts_only() {
560 let usage = serde_json::json!({
561 "prompt_tokens": 200,
562 "completion_tokens": 50
563 });
564
565 let telemetry = ProviderTelemetry::from_openai_usage(&usage, Some("req-abc"));
566
567 assert_eq!(telemetry.source, source::OPENAI_USAGE);
568 assert_eq!(telemetry.server_prompt_tokens, Some(200));
569 assert_eq!(telemetry.server_output_tokens, Some(50));
570 assert_eq!(telemetry.server_prompt_eval_ms, None);
571 assert_eq!(telemetry.request_id.as_deref(), Some("req-abc"));
572 }
573
574 #[test]
575 fn llamacpp_timings_promotes_source_and_fills_durations() {
576 let usage = serde_json::json!({
577 "prompt_tokens": 220,
578 "completion_tokens": 17,
579 "timings": {
580 "prompt_n": 200,
581 "prompt_ms": 145.4,
582 "predicted_n": 17,
583 "predicted_ms": 89.1,
584 }
585 });
586
587 let telemetry = ProviderTelemetry::from_openai_usage(&usage, None);
588
589 assert_eq!(telemetry.source, source::LLAMACPP_TIMINGS);
590 assert_eq!(telemetry.server_prompt_eval_ms, Some(145));
591 assert_eq!(telemetry.server_generation_ms, Some(89));
592 assert_eq!(telemetry.server_total_ms, Some(234));
593 assert_eq!(telemetry.server_prompt_tokens, Some(200));
594 assert_eq!(telemetry.server_output_tokens, Some(17));
595 assert!(!telemetry.is_empty());
596 }
597
598 #[test]
599 fn ps_entry_pulls_context_length_from_top_level_or_details() {
600 let entry = serde_json::json!({
601 "name": "devstral-small-2:24b",
602 "size": 4_700_000_000u64,
603 "size_vram": 4_500_000_000u64,
604 "expires_at": "2026-05-14T10:30:00Z",
605 "context_length": 32768
606 });
607 let model = OllamaPsModel::from_ps_entry(&entry).expect("ps entry parses");
608 assert_eq!(model.context_length, Some(32768));
609
610 let entry_nested = serde_json::json!({
611 "name": "devstral-small-2:24b",
612 "details": {"context_length": 16384}
613 });
614 let nested = OllamaPsModel::from_ps_entry(&entry_nested).expect("ps entry parses");
615 assert_eq!(nested.context_length, Some(16384));
616 }
617
618 #[test]
619 fn merge_ollama_ps_preserves_call_level_values() {
620 let mut telemetry = ProviderTelemetry::new(source::OLLAMA_CHAT);
621 telemetry.runtime_loaded_model = Some("real-model".to_string());
622 let ps = OllamaPsModel {
623 name: Some("alias-model".to_string()),
624 size_bytes: Some(1),
625 size_vram_bytes: Some(2),
626 expires_at: Some("forever".to_string()),
627 context_length: Some(8192),
628 };
629 telemetry.merge_ollama_ps(&ps);
630 assert_eq!(
631 telemetry.runtime_loaded_model.as_deref(),
632 Some("real-model")
633 );
634 assert_eq!(telemetry.runtime_memory_bytes, Some(1));
635 assert_eq!(telemetry.runtime_memory_vram_bytes, Some(2));
636 assert_eq!(
637 telemetry.runtime_keep_alive_until.as_deref(),
638 Some("forever")
639 );
640 assert_eq!(telemetry.runtime_context_length, Some(8192));
641 }
642
643 #[test]
644 fn as_vm_dict_returns_none_when_empty() {
645 let telemetry = ProviderTelemetry::default();
646 assert!(telemetry.is_empty());
647 assert!(telemetry.as_vm_dict().is_none());
648 }
649
650 #[test]
651 fn as_vm_dict_serializes_all_present_fields() {
652 let telemetry = ProviderTelemetry {
653 source: source::OLLAMA_CHAT.to_string(),
654 serving_base_url: Some("https://provider.example/v1".to_string()),
655 server_total_ms: Some(100),
656 client_wall_ms: Some(120),
657 runtime_loaded_model: Some("qwen".to_string()),
658 ..Default::default()
659 };
660 let value = telemetry.as_vm_dict().expect("dict present");
661 let dict = value.as_dict().expect("dict body");
662 assert_eq!(
663 dict.get("source").map(VmValue::display).as_deref(),
664 Some(source::OLLAMA_CHAT)
665 );
666 assert_eq!(
667 dict.get("serving_base_url")
668 .map(VmValue::display)
669 .as_deref(),
670 Some("https://provider.example/v1")
671 );
672 assert_eq!(
673 dict.get("server_total_ms").and_then(|v| match v {
674 VmValue::Int(n) => Some(*n),
675 _ => None,
676 }),
677 Some(100)
678 );
679 assert_eq!(
680 dict.get("client_wall_ms").and_then(|v| match v {
681 VmValue::Int(n) => Some(*n),
682 _ => None,
683 }),
684 Some(120)
685 );
686 }
687
688 #[test]
689 fn gateway_provider_metadata_is_preserved_without_schema_coupling() {
690 let response = serde_json::json!({
691 "model": "served-adapter",
692 "provider_metadata": {
693 "gateway": {
694 "routing": {
695 "resolvedProvider": "openai",
696 "modelAttemptCount": 1
697 },
698 "cost": "0.00003865"
699 }
700 }
701 });
702 let mut telemetry = ProviderTelemetry::default();
703 telemetry.capture_provider_metadata(&response);
704
705 assert_eq!(
706 telemetry
707 .provider_metadata
708 .as_ref()
709 .and_then(|metadata| metadata.pointer("/gateway/routing/resolvedProvider"))
710 .and_then(serde_json::Value::as_str),
711 Some("openai")
712 );
713 assert_eq!(telemetry.response_model.as_deref(), Some("served-adapter"));
714 assert!(!telemetry.is_empty());
715 let value = telemetry
716 .as_vm_dict()
717 .expect("metadata makes telemetry visible");
718 let dict = value.as_dict().expect("dict body");
719 assert!(dict.get("provider_metadata").is_some());
720 assert_eq!(
721 dict.get("response_model").map(VmValue::display).as_deref(),
722 Some("served-adapter")
723 );
724 }
725}