use crate::error::Result;
use crate::substrate::{ReadOpts, SubstrateRead};
use serde_json::Value;
pub const EVAL_RUN_RELATION: &str = "mg:eval_run";
pub const HARNESS_NS: &str = "agent:harness";
#[derive(Debug, Clone)]
pub struct EvalRun {
pub run_id: String,
pub passed: u64,
pub failed: u64,
pub summary: serde_json::Map<String, Value>,
pub recorded_ms: i64,
pub spend: Option<RunSpend>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RunSpend {
pub input_tokens: u64,
pub output_tokens: u64,
pub usd_micros: u64,
pub wall_ms: u64,
}
pub const COST_FIELDS: [&str; 5] = ["effects", "tokens", "usd", "wall_ms", "cost_per_pass"];
impl EvalRun {
pub fn field(&self, name: &str) -> Option<f64> {
self.summary.get(name).and_then(Value::as_f64)
}
pub fn total(&self) -> u64 {
self.passed.saturating_add(self.failed)
}
fn cost_key(&self, key: &str, from_spend: impl Fn(&RunSpend) -> u64) -> Option<u64> {
match self.summary.get(key) {
Some(v) => v.as_u64(),
None => self.spend.as_ref().map(from_spend),
}
}
pub fn effects(&self) -> Option<u64> {
self.summary.get("effects").and_then(Value::as_u64)
}
pub fn tokens(&self) -> Option<u64> {
let i = self.cost_key("input_tokens", |s| s.input_tokens)?;
let o = self.cost_key("output_tokens", |s| s.output_tokens)?;
Some(i.saturating_add(o))
}
pub fn usd_micros(&self) -> Option<u64> {
self.cost_key("usd_micros", |s| s.usd_micros)
}
pub fn wall_ms(&self) -> Option<u64> {
self.cost_key("wall_ms", |s| s.wall_ms)
}
}
pub fn eval_runs<S: SubstrateRead + ?Sized>(
sub: &S,
evalset_hash: &str,
since_ms: Option<i64>,
) -> Result<Vec<EvalRun>> {
let subject = format!("evalset:{evalset_hash}");
let facts = sub.grains_of_type(
crate::model::grain_type::FACT,
Some(HARNESS_NS),
ReadOpts { live_only: true, since_ms },
)?;
let spends = run_spends(sub)?;
let mut out: Vec<EvalRun> = facts
.iter()
.filter(|f| f.str_field("relation") == Some(EVAL_RUN_RELATION))
.filter(|f| f.str_field("subject") == Some(subject.as_str()))
.filter_map(|f| {
let obj = f.str_field("object")?;
let Ok(Value::Object(summary)) = serde_json::from_str::<Value>(obj) else {
return None;
};
let run_id = summary.get("run_id").and_then(Value::as_str)?.to_string();
Some(EvalRun {
spend: spends.get(&run_id).copied(),
run_id,
passed: summary.get("passed").and_then(Value::as_u64)?,
failed: summary.get("failed").and_then(Value::as_u64)?,
recorded_ms: f.created_at_ms,
summary,
})
})
.collect();
out.sort_by(|a, b| {
a.recorded_ms
.cmp(&b.recorded_ms)
.then_with(|| a.run_id.cmp(&b.run_id))
});
Ok(out)
}
fn run_spends<S: SubstrateRead + ?Sized>(sub: &S) -> Result<std::collections::BTreeMap<String, RunSpend>> {
let obs = match sub.grains_of_type(
crate::model::grain_type::OBSERVATION,
Some(HARNESS_NS),
ReadOpts { live_only: true, since_ms: None },
) {
Ok(rows) => rows,
Err(_) => return Ok(Default::default()),
};
let mut out = std::collections::BTreeMap::new();
for g in &obs {
if g.str_field("observation_kind") != Some("run_outcome") {
continue;
}
let Some(run_id) = g.str_field("run_id") else { continue };
let int = |k: &str| g.fields.get(k).and_then(Value::as_u64);
let (Some(i), Some(o), Some(u), Some(w)) = (
int("spent_input_tokens"),
int("spent_output_tokens"),
int("spent_usd_micros"),
int("spent_wall_ms"),
) else {
continue;
};
out.insert(
run_id.to_string(),
RunSpend { input_tokens: i, output_tokens: o, usd_micros: u, wall_ms: w },
);
}
Ok(out)
}
pub fn newest_eval_run<S: SubstrateRead + ?Sized>(
sub: &S,
evalset_hash: &str,
since_ms: Option<i64>,
) -> Result<Option<EvalRun>> {
Ok(eval_runs(sub, evalset_hash, since_ms)?.pop())
}
pub fn run_value(run: &EvalRun, field: &str) -> Option<f64> {
match field {
"failed" => Some(run.failed as f64),
"passed" => Some(run.passed as f64),
"total" => Some(run.total() as f64),
"error_rate" => match run.total() {
0 => None,
t => Some(run.failed as f64 / t as f64),
},
"effects" => run.effects().map(|n| n as f64),
"tokens" => run.tokens().map(|n| n as f64),
"usd" => run.usd_micros().map(|n| n as f64 / 1e6),
"wall_ms" => run.wall_ms().map(|n| n as f64),
"cost_per_pass" => match run.passed {
0 => None,
p => run.usd_micros().map(|n| n as f64 / 1e6 / p as f64),
},
other => run.field(other),
}
}
pub fn eval_run_by_id<S: SubstrateRead + ?Sized>(
sub: &S,
evalset_hash: &str,
run_id: &str,
) -> Result<Option<EvalRun>> {
Ok(eval_runs(sub, evalset_hash, None)?
.into_iter()
.find(|r| r.run_id == run_id))
}
pub fn parse_evalset_metric(metric: &str) -> Option<(&str, &str)> {
let rest = metric.strip_prefix("evalset:")?;
let (hash, field) = rest.rsplit_once(':')?;
if hash.is_empty() || field.is_empty() || hash.contains(':') {
return None;
}
Some((hash, field))
}
#[cfg(test)]
mod tests {
use super::*;
fn run(summary: serde_json::Value, spend: Option<RunSpend>) -> EvalRun {
let summary = summary.as_object().unwrap().clone();
EvalRun {
run_id: "eval-1".into(),
passed: summary.get("passed").and_then(Value::as_u64).unwrap_or(0),
failed: summary.get("failed").and_then(Value::as_u64).unwrap_or(0),
summary,
recorded_ms: 0,
spend,
}
}
#[test]
fn cost_fields_are_integers_or_not_measurable() {
let ok = run(
serde_json::json!({"passed": 4, "failed": 1, "effects": 12, "input_tokens": 1000,
"output_tokens": 200, "usd_micros": 2_000_000, "wall_ms": 340}),
None,
);
assert_eq!(run_value(&ok, "effects"), Some(12.0));
assert_eq!(run_value(&ok, "tokens"), Some(1200.0));
assert_eq!(run_value(&ok, "usd"), Some(2.0));
assert_eq!(run_value(&ok, "wall_ms"), Some(340.0));
assert_eq!(run_value(&ok, "cost_per_pass"), Some(0.5));
for bad in [
serde_json::json!({"passed": 4, "failed": 1, "input_tokens": "1000", "output_tokens": 200}),
serde_json::json!({"passed": 4, "failed": 1, "input_tokens": 1000.5, "output_tokens": 200}),
serde_json::json!({"passed": 4, "failed": 1, "input_tokens": -1, "output_tokens": 200}),
serde_json::json!({"passed": 4, "failed": 1, "output_tokens": 200}),
serde_json::json!({"passed": 4, "failed": 1}),
] {
let r = run(bad.clone(), None);
assert_eq!(run_value(&r, "tokens"), None, "{bad}");
assert_eq!(run_value(&r, "cost_per_pass"), None, "{bad}");
}
let none = run(serde_json::json!({"passed": 0, "failed": 5, "usd_micros": 2_000_000}), None);
assert_eq!(run_value(&none, "usd"), Some(2.0));
assert_eq!(run_value(&none, "cost_per_pass"), None);
let r = run(serde_json::json!({"passed": 4, "failed": 1, "wall_ms": "fast"}), None);
assert_eq!(run_value(&r, "passed"), Some(4.0));
assert_eq!(run_value(&r, "wall_ms"), None);
}
#[test]
fn a_missing_cost_key_falls_back_to_the_runtime_spend_and_a_present_one_does_not() {
let spend = RunSpend { input_tokens: 700, output_tokens: 300, usd_micros: 4_000_000, wall_ms: 9_000 };
let r = run(serde_json::json!({"passed": 2, "failed": 0}), Some(spend));
assert_eq!(run_value(&r, "tokens"), Some(1000.0));
assert_eq!(run_value(&r, "usd"), Some(4.0));
assert_eq!(run_value(&r, "wall_ms"), Some(9000.0));
assert_eq!(run_value(&r, "cost_per_pass"), Some(2.0));
assert_eq!(run_value(&r, "effects"), None, "the runtime records supersteps, not effects");
let wrong = run(serde_json::json!({"passed": 2, "failed": 0, "input_tokens": "700"}), Some(spend));
assert_eq!(run_value(&wrong, "tokens"), None, "a malformed key is malformed, not missing");
let own = run(serde_json::json!({"passed": 2, "failed": 0, "input_tokens": 1, "output_tokens": 1}), Some(spend));
assert_eq!(run_value(&own, "tokens"), Some(2.0), "the summary's own figure wins");
}
#[test]
fn metric_strings_parse_into_hash_and_field() {
assert_eq!(
parse_evalset_metric("evalset:abc123:category_accuracy"),
Some(("abc123", "category_accuracy"))
);
assert_eq!(parse_evalset_metric("evalset:abc123:failed"), Some(("abc123", "failed")));
assert_eq!(parse_evalset_metric("tool_error_recurrence"), None);
assert_eq!(parse_evalset_metric("evalset:abc123"), None);
assert_eq!(parse_evalset_metric("evalset::field"), None);
assert_eq!(parse_evalset_metric("evalset:abc123:"), None);
assert_eq!(parse_evalset_metric("evalset:a:b:c"), None);
}
}