systemprompt_analytics/
resource_metrics.rs1use chrono::{DateTime, Utc};
8use serde::Serialize;
9use std::collections::{BTreeMap, BTreeSet};
10use systemprompt_identifiers::{AiRequestId, ResourceInvocationId, SessionId, UserId};
11
12#[derive(Debug, Clone)]
13pub struct ResourceFact {
14 pub invocation_id: ResourceInvocationId,
15 pub user_id: UserId,
16 pub session_id: SessionId,
17 pub invoked_at: DateTime<Utc>,
18 pub request_id: Option<AiRequestId>,
19 pub input_tokens: Option<i64>,
20 pub output_tokens: Option<i64>,
21 pub cache_read_tokens: Option<i64>,
22 pub cache_creation_tokens: Option<i64>,
23 pub cost_microdollars: Option<i64>,
24 pub latency_ms: Option<i64>,
25 pub failed: bool,
26 pub quality_score: Option<f64>,
27 pub successful: Option<bool>,
28 pub revision_verified: bool,
29}
30
31#[derive(Debug, Default, Clone, Copy, Serialize)]
32pub struct ResourceMetrics {
33 pub invocations: usize,
34 pub users: usize,
35 pub conversations: usize,
36 pub requests: usize,
37 pub measured_requests: usize,
38 pub priced_requests: usize,
39 pub failed_requests: usize,
40 pub verified_invocations: usize,
41 pub assessed_conversations: usize,
42 pub successful_conversations: usize,
43 pub input_tokens: i128,
44 pub output_tokens: i128,
45 pub cache_read_tokens: i128,
46 pub cache_creation_tokens: i128,
47 pub related_cost_microdollars: Option<i128>,
48 pub average_tokens_per_measured_request: Option<f64>,
49 pub average_latency_ms: Option<f64>,
50 pub average_quality_score: Option<f64>,
51 pub last_used_at: Option<DateTime<Utc>>,
52}
53
54pub fn aggregate<'a>(facts: impl IntoIterator<Item = &'a ResourceFact> + 'a) -> ResourceMetrics {
55 let mut invocations = BTreeSet::new();
56 let mut verified = BTreeSet::new();
57 let mut users = BTreeSet::new();
58 let mut sessions = BTreeSet::new();
59 let mut requests = BTreeMap::new();
60 let mut assessments = BTreeMap::new();
61 let mut last_used = None;
62 for fact in facts {
63 invocations.insert((&fact.user_id, &fact.invocation_id));
64 if fact.revision_verified {
65 verified.insert((&fact.user_id, &fact.invocation_id));
66 }
67 users.insert(&fact.user_id);
68 sessions.insert((&fact.user_id, &fact.session_id));
69 last_used = Some(last_used.map_or(fact.invoked_at, |last: DateTime<Utc>| {
70 last.max(fact.invoked_at)
71 }));
72 if fact.quality_score.is_some() || fact.successful.is_some() {
73 assessments.insert(
74 (&fact.user_id, &fact.session_id),
75 (fact.quality_score, fact.successful),
76 );
77 }
78 if let Some(request) = &fact.request_id {
79 requests.entry((&fact.user_id, request)).or_insert(fact);
80 }
81 }
82 let mut result = ResourceMetrics {
83 invocations: invocations.len(),
84 users: users.len(),
85 conversations: sessions.len(),
86 requests: requests.len(),
87 verified_invocations: verified.len(),
88 assessed_conversations: assessments.len(),
89 last_used_at: last_used,
90 ..ResourceMetrics::default()
91 };
92 aggregate_requests(&mut result, requests.into_values());
93 result.average_quality_score = mean(
94 &assessments
95 .values()
96 .filter_map(|(score, _)| *score)
97 .collect::<Vec<_>>(),
98 );
99 result.successful_conversations = assessments
100 .values()
101 .filter(|(_, success)| *success == Some(true))
102 .count();
103 result
104}
105
106fn aggregate_requests<'a>(
107 result: &mut ResourceMetrics,
108 facts: impl IntoIterator<Item = &'a ResourceFact> + 'a,
109) {
110 let mut token_sum = 0i128;
111 let mut latencies = Vec::new();
112 let mut cost = 0i128;
113 for fact in facts {
114 result.failed_requests += usize::from(fact.failed);
115 if let Some(amount) = fact.cost_microdollars.filter(|value| *value >= 0) {
116 result.priced_requests += 1;
117 cost += i128::from(amount);
118 }
119 if let (Some(input), Some(output)) = (
120 fact.input_tokens.filter(|value| *value >= 0),
121 fact.output_tokens.filter(|value| *value >= 0),
122 ) {
123 result.measured_requests += 1;
124 let cache_read = i128::from(fact.cache_read_tokens.unwrap_or(0).max(0));
125 let cache_write = i128::from(fact.cache_creation_tokens.unwrap_or(0).max(0));
126 result.input_tokens += i128::from(input);
127 result.output_tokens += i128::from(output);
128 result.cache_read_tokens += cache_read;
129 result.cache_creation_tokens += cache_write;
130 token_sum += i128::from(input) + i128::from(output) + cache_read + cache_write;
131 }
132 if let Some(latency) = fact.latency_ms {
133 latencies.push(latency as f64);
134 }
135 }
136 result.related_cost_microdollars = (result.priced_requests > 0).then_some(cost);
137 result.average_tokens_per_measured_request =
138 (result.measured_requests > 0).then(|| token_sum as f64 / result.measured_requests as f64);
139 result.average_latency_ms = mean(&latencies);
140}
141
142fn mean(values: &[f64]) -> Option<f64> {
143 (!values.is_empty()).then(|| values.iter().sum::<f64>() / values.len() as f64)
144}