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 self.provider_metadata = Some(metadata.clone());
330 }
331 }
332
333 pub fn merge_ollama_ps(&mut self, ps: &OllamaPsModel) {
337 if self.runtime_loaded_model.is_none() {
338 self.runtime_loaded_model = ps.name.clone();
339 }
340 if self.runtime_context_length.is_none() {
341 self.runtime_context_length = ps.context_length;
342 }
343 if self.runtime_memory_bytes.is_none() {
344 self.runtime_memory_bytes = ps.size_bytes;
345 }
346 if self.runtime_memory_vram_bytes.is_none() {
347 self.runtime_memory_vram_bytes = ps.size_vram_bytes;
348 }
349 if self.runtime_keep_alive_until.is_none() {
350 self.runtime_keep_alive_until = ps.expires_at.clone();
351 }
352 }
353
354 pub fn as_vm_dict(&self) -> Option<VmValue> {
358 if self.is_empty() {
359 return None;
360 }
361 let mut dict: crate::value::DictMap = crate::value::DictMap::new();
362 if !self.source.is_empty() {
363 dict.put_str("source", self.source.as_str());
364 }
365 if let Some(ref serving_base_url) = self.serving_base_url {
366 dict.put_str("serving_base_url", serving_base_url.as_str());
367 }
368 insert_opt_u64(&mut dict, "server_total_ms", self.server_total_ms);
369 insert_opt_u64(&mut dict, "server_load_ms", self.server_load_ms);
370 insert_opt_u64(
371 &mut dict,
372 "server_prompt_eval_ms",
373 self.server_prompt_eval_ms,
374 );
375 insert_opt_u64(&mut dict, "server_generation_ms", self.server_generation_ms);
376 insert_opt_i64(&mut dict, "server_prompt_tokens", self.server_prompt_tokens);
377 insert_opt_i64(&mut dict, "server_output_tokens", self.server_output_tokens);
378 insert_opt_u64(&mut dict, "client_wall_ms", self.client_wall_ms);
379 insert_opt_u64(
380 &mut dict,
381 "runtime_context_length",
382 self.runtime_context_length,
383 );
384 if let Some(ref model) = self.runtime_loaded_model {
385 dict.put_str("runtime_loaded_model", model.as_str());
386 }
387 if let Some(ref model) = self.response_model {
388 dict.put_str("response_model", model.as_str());
389 }
390 insert_opt_u64(&mut dict, "runtime_memory_bytes", self.runtime_memory_bytes);
391 insert_opt_u64(
392 &mut dict,
393 "runtime_memory_vram_bytes",
394 self.runtime_memory_vram_bytes,
395 );
396 if let Some(ref expires) = self.runtime_keep_alive_until {
397 dict.put_str("runtime_keep_alive_until", expires.as_str());
398 }
399 if let Some(ref request_id) = self.request_id {
400 dict.put_str("request_id", request_id.as_str());
401 }
402 if let Some(ref provider_metadata) = self.provider_metadata {
403 dict.insert(
404 crate::value::intern_key("provider_metadata"),
405 crate::stdlib::json_to_vm_value(provider_metadata),
406 );
407 }
408 Some(VmValue::dict(dict))
409 }
410}
411
412#[derive(Clone, Debug, Default, PartialEq, Eq)]
415pub struct OllamaPsModel {
416 pub name: Option<String>,
417 pub size_bytes: Option<u64>,
418 pub size_vram_bytes: Option<u64>,
419 pub expires_at: Option<String>,
420 pub context_length: Option<u64>,
421}
422
423impl OllamaPsModel {
424 pub fn from_ps_entry(entry: &serde_json::Value) -> Option<Self> {
428 let name = entry
429 .get("name")
430 .and_then(serde_json::Value::as_str)
431 .or_else(|| entry.get("model").and_then(serde_json::Value::as_str))
432 .map(str::to_string);
433 let context_length = entry
434 .get("context_length")
435 .and_then(serde_json::Value::as_u64)
436 .or_else(|| {
437 entry
438 .get("details")
439 .and_then(|details| details.get("context_length"))
440 .and_then(serde_json::Value::as_u64)
441 });
442 Some(Self {
443 name,
444 size_bytes: entry.get("size").and_then(serde_json::Value::as_u64),
445 size_vram_bytes: entry.get("size_vram").and_then(serde_json::Value::as_u64),
446 expires_at: entry
447 .get("expires_at")
448 .and_then(serde_json::Value::as_str)
449 .map(str::to_string),
450 context_length,
451 })
452 }
453}
454
455fn ns_field(frame: &serde_json::Value, key: &str) -> Option<u64> {
456 frame
457 .get(key)
458 .and_then(serde_json::Value::as_u64)
459 .map(ProviderTelemetry::ns_to_ms)
460}
461
462fn ms_or_round(value: Option<&serde_json::Value>) -> Option<u64> {
463 let value = value?;
464 if let Some(n) = value.as_u64() {
465 return Some(n);
466 }
467 value.as_f64().map(|n| n.round().max(0.0) as u64)
468}
469
470fn insert_opt_u64(dict: &mut crate::value::DictMap, key: &str, value: Option<u64>) {
471 if let Some(value) = value {
472 dict.insert(crate::value::intern_key(key), VmValue::Int(value as i64));
473 }
474}
475
476fn insert_opt_i64(dict: &mut crate::value::DictMap, key: &str, value: Option<i64>) {
477 if let Some(value) = value {
478 dict.insert(crate::value::intern_key(key), VmValue::Int(value));
479 }
480}
481
482#[cfg(test)]
483mod tests {
484 use super::*;
485
486 #[test]
487 fn ollama_done_frame_extracts_full_breakdown() {
488 let frame = serde_json::json!({
489 "model": "devstral-small-2:24b",
490 "total_duration": 7_400_000_000u64,
491 "load_duration": 400_000_000u64,
492 "prompt_eval_duration": 1_200_000_000u64,
493 "eval_duration": 5_800_000_000u64,
494 "prompt_eval_count": 1024,
495 "eval_count": 64
496 });
497
498 let telemetry = ProviderTelemetry::from_ollama_done(&frame, source::OLLAMA_CHAT);
499
500 assert_eq!(telemetry.source, source::OLLAMA_CHAT);
501 assert_eq!(telemetry.server_total_ms, Some(7400));
502 assert_eq!(telemetry.server_load_ms, Some(400));
503 assert_eq!(telemetry.server_prompt_eval_ms, Some(1200));
504 assert_eq!(telemetry.server_generation_ms, Some(5800));
505 assert_eq!(telemetry.server_prompt_tokens, Some(1024));
506 assert_eq!(telemetry.server_output_tokens, Some(64));
507 assert_eq!(
508 telemetry.runtime_loaded_model.as_deref(),
509 Some("devstral-small-2:24b")
510 );
511 assert!(!telemetry.is_empty());
512 }
513
514 #[test]
515 fn ollama_done_frame_leaves_missing_fields_as_none() {
516 let frame = serde_json::json!({
517 "model": "devstral-small-2:24b",
518 });
520
521 let telemetry = ProviderTelemetry::from_ollama_done(&frame, source::OLLAMA_CHAT);
522
523 assert_eq!(telemetry.server_total_ms, None);
524 assert_eq!(telemetry.server_load_ms, None);
525 assert_eq!(telemetry.server_prompt_eval_ms, None);
526 assert_eq!(telemetry.server_generation_ms, None);
527 assert_eq!(telemetry.server_prompt_tokens, None);
528 assert_eq!(telemetry.server_output_tokens, None);
529 }
530
531 #[test]
532 fn openai_usage_partial_extracts_counts_only() {
533 let usage = serde_json::json!({
534 "prompt_tokens": 200,
535 "completion_tokens": 50
536 });
537
538 let telemetry = ProviderTelemetry::from_openai_usage(&usage, Some("req-abc"));
539
540 assert_eq!(telemetry.source, source::OPENAI_USAGE);
541 assert_eq!(telemetry.server_prompt_tokens, Some(200));
542 assert_eq!(telemetry.server_output_tokens, Some(50));
543 assert_eq!(telemetry.server_prompt_eval_ms, None);
544 assert_eq!(telemetry.request_id.as_deref(), Some("req-abc"));
545 }
546
547 #[test]
548 fn llamacpp_timings_promotes_source_and_fills_durations() {
549 let usage = serde_json::json!({
550 "prompt_tokens": 220,
551 "completion_tokens": 17,
552 "timings": {
553 "prompt_n": 200,
554 "prompt_ms": 145.4,
555 "predicted_n": 17,
556 "predicted_ms": 89.1,
557 }
558 });
559
560 let telemetry = ProviderTelemetry::from_openai_usage(&usage, None);
561
562 assert_eq!(telemetry.source, source::LLAMACPP_TIMINGS);
563 assert_eq!(telemetry.server_prompt_eval_ms, Some(145));
564 assert_eq!(telemetry.server_generation_ms, Some(89));
565 assert_eq!(telemetry.server_total_ms, Some(234));
566 assert_eq!(telemetry.server_prompt_tokens, Some(200));
567 assert_eq!(telemetry.server_output_tokens, Some(17));
568 assert!(!telemetry.is_empty());
569 }
570
571 #[test]
572 fn ps_entry_pulls_context_length_from_top_level_or_details() {
573 let entry = serde_json::json!({
574 "name": "devstral-small-2:24b",
575 "size": 4_700_000_000u64,
576 "size_vram": 4_500_000_000u64,
577 "expires_at": "2026-05-14T10:30:00Z",
578 "context_length": 32768
579 });
580 let model = OllamaPsModel::from_ps_entry(&entry).expect("ps entry parses");
581 assert_eq!(model.context_length, Some(32768));
582
583 let entry_nested = serde_json::json!({
584 "name": "devstral-small-2:24b",
585 "details": {"context_length": 16384}
586 });
587 let nested = OllamaPsModel::from_ps_entry(&entry_nested).expect("ps entry parses");
588 assert_eq!(nested.context_length, Some(16384));
589 }
590
591 #[test]
592 fn merge_ollama_ps_preserves_call_level_values() {
593 let mut telemetry = ProviderTelemetry::new(source::OLLAMA_CHAT);
594 telemetry.runtime_loaded_model = Some("real-model".to_string());
595 let ps = OllamaPsModel {
596 name: Some("alias-model".to_string()),
597 size_bytes: Some(1),
598 size_vram_bytes: Some(2),
599 expires_at: Some("forever".to_string()),
600 context_length: Some(8192),
601 };
602 telemetry.merge_ollama_ps(&ps);
603 assert_eq!(
604 telemetry.runtime_loaded_model.as_deref(),
605 Some("real-model")
606 );
607 assert_eq!(telemetry.runtime_memory_bytes, Some(1));
608 assert_eq!(telemetry.runtime_memory_vram_bytes, Some(2));
609 assert_eq!(
610 telemetry.runtime_keep_alive_until.as_deref(),
611 Some("forever")
612 );
613 assert_eq!(telemetry.runtime_context_length, Some(8192));
614 }
615
616 #[test]
617 fn as_vm_dict_returns_none_when_empty() {
618 let telemetry = ProviderTelemetry::default();
619 assert!(telemetry.is_empty());
620 assert!(telemetry.as_vm_dict().is_none());
621 }
622
623 #[test]
624 fn as_vm_dict_serializes_all_present_fields() {
625 let telemetry = ProviderTelemetry {
626 source: source::OLLAMA_CHAT.to_string(),
627 serving_base_url: Some("https://provider.example/v1".to_string()),
628 server_total_ms: Some(100),
629 client_wall_ms: Some(120),
630 runtime_loaded_model: Some("qwen".to_string()),
631 ..Default::default()
632 };
633 let value = telemetry.as_vm_dict().expect("dict present");
634 let dict = value.as_dict().expect("dict body");
635 assert_eq!(
636 dict.get("source").map(VmValue::display).as_deref(),
637 Some(source::OLLAMA_CHAT)
638 );
639 assert_eq!(
640 dict.get("serving_base_url")
641 .map(VmValue::display)
642 .as_deref(),
643 Some("https://provider.example/v1")
644 );
645 assert_eq!(
646 dict.get("server_total_ms").and_then(|v| match v {
647 VmValue::Int(n) => Some(*n),
648 _ => None,
649 }),
650 Some(100)
651 );
652 assert_eq!(
653 dict.get("client_wall_ms").and_then(|v| match v {
654 VmValue::Int(n) => Some(*n),
655 _ => None,
656 }),
657 Some(120)
658 );
659 }
660
661 #[test]
662 fn gateway_provider_metadata_is_preserved_without_schema_coupling() {
663 let response = serde_json::json!({
664 "model": "served-adapter",
665 "provider_metadata": {
666 "gateway": {
667 "routing": {
668 "resolvedProvider": "openai",
669 "modelAttemptCount": 1
670 },
671 "cost": "0.00003865"
672 }
673 }
674 });
675 let mut telemetry = ProviderTelemetry::default();
676 telemetry.capture_provider_metadata(&response);
677
678 assert_eq!(
679 telemetry
680 .provider_metadata
681 .as_ref()
682 .and_then(|metadata| metadata.pointer("/gateway/routing/resolvedProvider"))
683 .and_then(serde_json::Value::as_str),
684 Some("openai")
685 );
686 assert_eq!(telemetry.response_model.as_deref(), Some("served-adapter"));
687 assert!(!telemetry.is_empty());
688 let value = telemetry
689 .as_vm_dict()
690 .expect("metadata makes telemetry visible");
691 let dict = value.as_dict().expect("dict body");
692 assert!(dict.get("provider_metadata").is_some());
693 assert_eq!(
694 dict.get("response_model").map(VmValue::display).as_deref(),
695 Some("served-adapter")
696 );
697 }
698}