Skip to main content

systemprompt_analytics/
resource_metrics.rs

1//! Resource metrics count each request and assessed conversation once within
2//! a cohort. Related conversation spend remains non-additive across cohorts.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use 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}