Skip to main content

areev_loop/analyzers/
run_outcome.rs

1//! Run outcomes (T0) — whole-run health for `areev run` workflows (§8 Wave
2//! 4). The datasource is the compact `run_outcome` Observation the driver
3//! writes at every terminal run (outcome label + spent figures + plan
4//! hash). Two signals, both advisory (Flag — what to do about a failing or
5//! expensive workflow is a human/host decision, never auto-applied):
6//!
7//! - **Failure clusters per workflow**: the same plan hash finishing
8//!   `failed`/`stalled`/`budget_exhausted` repeatedly.
9//! - **Cost attribution**: aggregate USD spend per workflow — the run-side
10//!   `budget_pressure` signal (assembly-budget pressure has its own
11//!   analyzer; this one watches what runs actually spend).
12
13use crate::analyzer::{AnalyzeCtx, Analyzer};
14use crate::error::Result;
15use crate::manifest::*;
16use crate::model::{ActionKind, Severity};
17use crate::recommendation::{Proposal, RecDraft, Summary};
18use serde_json::{json, Map};
19use std::collections::BTreeMap;
20
21/// The namespace `areev run` journals into.
22const HARNESS_NS: &str = "agent:harness";
23
24pub struct RunOutcome {
25    manifest: AnalyzerManifest,
26}
27
28impl RunOutcome {
29    pub fn new() -> Self {
30        RunOutcome {
31            manifest: AnalyzerManifest {
32                id: "loop.run_outcome/1".into(),
33                title: "Run outcomes".into(),
34                description:
35                    "Flags workflows whose runs keep failing, stalling, or exhausting \
36                     budgets, and attributes run spend per workflow."
37                        .into(),
38                tier: Tier::T0,
39                cadence: CadenceClass::Slow,
40                requires: vec![],
41                target_classes: vec![TargetClass::Host],
42                auto_apply: AutoApplyClass::Never,
43                trust_class: TrustClass::Builtin,
44                params: vec![
45                    ParamSpec::Int {
46                        name: "min_runs".into(),
47                        default: 3,
48                        min: 1,
49                        max: 1_000_000,
50                        description: "Minimum terminal runs of one workflow before its \
51                                      failure rate is meaningful."
52                            .into(),
53                    },
54                    ParamSpec::Float {
55                        name: "min_failure_ratio".into(),
56                        default: 0.5,
57                        min: 0.0,
58                        max: 1.0,
59                        description: "Non-completed fraction at or above which the \
60                                      workflow is flagged."
61                            .into(),
62                    },
63                    ParamSpec::Int {
64                        name: "min_usd_micros".into(),
65                        default: 5_000_000, // $5
66                        min: 0,
67                        max: i64::MAX,
68                        description: "Aggregate spend (micro-USD) at or above which a \
69                                      workflow's cost is surfaced."
70                            .into(),
71                    },
72                ],
73                default_on: true,
74            },
75        }
76    }
77}
78
79impl Default for RunOutcome {
80    fn default() -> Self {
81        Self::new()
82    }
83}
84
85#[derive(Default)]
86struct PlanStats {
87    runs: i64,
88    failed: i64,
89    usd_micros: i64,
90    last_detail: String,
91}
92
93impl Analyzer for RunOutcome {
94    fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
95        let obs = match ctx.grains_in("observation", HARNESS_NS) {
96            Ok(rows) => rows,
97            // No harness namespace / no read grant: no runs to analyze —
98            // degrade to nothing, never fabricate.
99            Err(_) => return Ok(Vec::new()),
100        };
101        let min_runs = ctx.params().get_int("min_runs");
102        let min_failure_ratio = ctx.params().get_float("min_failure_ratio");
103        let min_usd = ctx.params().get_int("min_usd_micros");
104
105        let mut by_plan: BTreeMap<String, PlanStats> = BTreeMap::new();
106        let mut evidence: BTreeMap<String, Vec<String>> = BTreeMap::new();
107        for g in &obs {
108            if g.str_field("observation_kind") != Some("run_outcome") {
109                continue;
110            }
111            let Some(plan) = g.str_field("plan_hash") else { continue };
112            let outcome = g.str_field("object").unwrap_or_default();
113            let entry = by_plan.entry(plan.to_string()).or_default();
114            entry.runs += 1;
115            if outcome != "completed" && outcome != "canceled" {
116                entry.failed += 1;
117                if let Some(d) = g.str_field("outcome_detail") {
118                    entry.last_detail = d.to_string();
119                } else {
120                    entry.last_detail = outcome.to_string();
121                }
122            }
123            entry.usd_micros += g
124                .fields
125                .get("spent_usd_micros")
126                .and_then(serde_json::Value::as_i64)
127                .unwrap_or(0);
128            evidence.entry(plan.to_string()).or_default().push(g.hash.clone());
129        }
130
131        let mut out = Vec::new();
132        for (plan, s) in &by_plan {
133            let short = &plan[..plan.len().min(12)];
134            if s.runs >= min_runs
135                && (s.failed as f64 / s.runs as f64) >= min_failure_ratio
136            {
137                let mut args = Map::new();
138                args.insert("workflow".into(), json!(short));
139                args.insert("failed".into(), json!(s.failed));
140                args.insert("runs".into(), json!(s.runs));
141                args.insert(
142                    "rate".into(),
143                    json!(((s.failed as f64 / s.runs as f64) * 100.0).round() as i64),
144                );
145                args.insert("last_error".into(), json!(s.last_detail));
146                let mut data = Map::new();
147                data.insert("plan_hash".into(), json!(plan));
148                data.insert("failed".into(), json!(s.failed));
149                data.insert("runs".into(), json!(s.runs));
150                out.push(
151                    RecDraft::new(
152                        format!("host:workflow/{plan}"),
153                        ActionKind::Flag,
154                        Summary::new("run.failures", args),
155                        Proposal::Data { data },
156                    )
157                    .severity(Severity::High)
158                    .evidence(evidence.get(plan).cloned().unwrap_or_default()),
159                );
160            }
161            if min_usd > 0 && s.usd_micros >= min_usd && s.runs > 0 {
162                let mut args = Map::new();
163                args.insert("workflow".into(), json!(short));
164                args.insert("runs".into(), json!(s.runs));
165                args.insert(
166                    "usd".into(),
167                    json!(format!("{:.2}", s.usd_micros as f64 / 1e6)),
168                );
169                args.insert(
170                    "avg_usd".into(),
171                    json!(format!("{:.2}", s.usd_micros as f64 / 1e6 / s.runs as f64)),
172                );
173                let mut data = Map::new();
174                data.insert("plan_hash".into(), json!(plan));
175                data.insert("usd_micros".into(), json!(s.usd_micros));
176                out.push(
177                    RecDraft::new(
178                        format!("host:workflow-cost/{plan}"),
179                        ActionKind::Flag,
180                        Summary::new("run.cost", args),
181                        Proposal::Data { data },
182                    )
183                    .severity(Severity::Medium)
184                    .evidence(evidence.get(plan).cloned().unwrap_or_default()),
185                );
186            }
187        }
188        Ok(out)
189    }
190
191    fn manifest(&self) -> &AnalyzerManifest {
192        &self.manifest
193    }
194}
195
196#[cfg(test)]
197mod tests {
198    use super::*;
199    use crate::testkit::TestSubstrate;
200
201    fn outcome(sub: &mut TestSubstrate, plan: &str, outcome: &str, usd: i64) {
202        sub.put_observation(
203            HARNESS_NS,
204            &[
205                ("observation_kind", json!("run_outcome")),
206                ("plan_hash", json!(plan)),
207                ("object", json!(outcome)),
208                ("spent_usd_micros", json!(usd)),
209                ("outcome_detail", json!("greet: ExecutorError: down")),
210            ],
211        );
212    }
213
214    #[test]
215    fn flags_failing_workflow() {
216        let mut sub = TestSubstrate::new();
217        for _ in 0..2 {
218            outcome(&mut sub, "abcd1234", "failed", 0);
219        }
220        outcome(&mut sub, "abcd1234", "completed", 0);
221        let drafts = sub.analyze(&RunOutcome::new(), 10_000);
222        assert_eq!(drafts.len(), 1, "{drafts:?}");
223        let text = drafts[0].summary.render();
224        assert!(text.contains("failed 2/3"), "{text}");
225        assert_eq!(drafts[0].action_kind, ActionKind::Flag);
226    }
227
228    #[test]
229    fn healthy_workflow_stays_quiet() {
230        let mut sub = TestSubstrate::new();
231        for _ in 0..5 {
232            outcome(&mut sub, "abcd1234", "completed", 100);
233        }
234        assert!(sub.analyze(&RunOutcome::new(), 10_000).is_empty());
235    }
236
237    #[test]
238    fn cost_attribution_surfaces_expensive_workflows() {
239        let mut sub = TestSubstrate::new();
240        for _ in 0..4 {
241            outcome(&mut sub, "eeff5566", "completed", 2_000_000); // $2 each
242        }
243        let drafts = sub.analyze(&RunOutcome::new(), 10_000);
244        assert_eq!(drafts.len(), 1);
245        let text = drafts[0].summary.render();
246        assert!(text.contains("8.00"), "aggregate spend rendered: {text}");
247    }
248}