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, attributes run spend per workflow, and surfaces those \
37                     whose transcripts keep outgrowing the model's window."
38                        .into(),
39                tier: Tier::T0,
40                cadence: CadenceClass::Slow,
41                requires: vec![],
42                target_classes: vec![TargetClass::Host],
43                auto_apply: AutoApplyClass::Never,
44                trust_class: TrustClass::Builtin,
45                params: vec![
46                    ParamSpec::Int {
47                        name: "min_runs".into(),
48                        default: 3,
49                        min: 1,
50                        max: 1_000_000,
51                        description: "Minimum terminal runs of one workflow before its \
52                                      failure rate is meaningful."
53                            .into(),
54                    },
55                    ParamSpec::Float {
56                        name: "min_failure_ratio".into(),
57                        default: 0.5,
58                        min: 0.0,
59                        max: 1.0,
60                        description: "Non-completed fraction at or above which the \
61                                      workflow is flagged."
62                            .into(),
63                    },
64                    ParamSpec::Int {
65                        name: "min_folds_per_run".into(),
66                        // Folding once in a while is the mechanism working. A
67                        // workflow averaging one fold PER RUN is one whose
68                        // nodes no longer fit, which is a design signal.
69                        default: 1,
70                        min: 1,
71                        max: 1_000_000,
72                        description: "Average transcript folds per run at or above which \
73                                      a workflow's context pressure is surfaced."
74                            .into(),
75                    },
76                    ParamSpec::Int {
77                        name: "min_usd_micros".into(),
78                        default: 5_000_000, // $5
79                        min: 0,
80                        max: i64::MAX,
81                        description: "Aggregate spend (micro-USD) at or above which a \
82                                      workflow's cost is surfaced."
83                            .into(),
84                    },
85                ],
86                default_on: true,
87            },
88        }
89    }
90}
91
92impl Default for RunOutcome {
93    fn default() -> Self {
94        Self::new()
95    }
96}
97
98#[derive(Default)]
99struct PlanStats {
100    runs: i64,
101    failed: i64,
102    usd_micros: i64,
103    folds: i64,
104    last_detail: String,
105}
106
107impl Analyzer for RunOutcome {
108    fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
109        let obs = match ctx.grains_in("observation", HARNESS_NS) {
110            Ok(rows) => rows,
111            // No harness namespace / no read grant: no runs to analyze —
112            // degrade to nothing, never fabricate.
113            Err(_) => return Ok(Vec::new()),
114        };
115        let min_runs = ctx.params().get_int("min_runs");
116        let min_failure_ratio = ctx.params().get_float("min_failure_ratio");
117        let min_usd = ctx.params().get_int("min_usd_micros");
118        let min_folds_per_run = ctx.params().get_int("min_folds_per_run");
119
120        let mut by_plan: BTreeMap<String, PlanStats> = BTreeMap::new();
121        let mut evidence: BTreeMap<String, Vec<String>> = BTreeMap::new();
122        for g in &obs {
123            if g.str_field("observation_kind") != Some("run_outcome") {
124                continue;
125            }
126            let Some(plan) = g.str_field("plan_hash") else { continue };
127            let outcome = g.str_field("object").unwrap_or_default();
128            let entry = by_plan.entry(plan.to_string()).or_default();
129            entry.runs += 1;
130            if outcome != "completed" && outcome != "canceled" {
131                entry.failed += 1;
132                if let Some(d) = g.str_field("outcome_detail") {
133                    entry.last_detail = d.to_string();
134                } else {
135                    entry.last_detail = outcome.to_string();
136                }
137            }
138            entry.usd_micros += g
139                .fields
140                .get("spent_usd_micros")
141                .and_then(serde_json::Value::as_i64)
142                .unwrap_or(0);
143            // Absent on every run that never folded, which is most of them.
144            entry.folds +=
145                g.fields.get("folds").and_then(serde_json::Value::as_i64).unwrap_or(0);
146            evidence.entry(plan.to_string()).or_default().push(g.hash.clone());
147        }
148
149        let mut out = Vec::new();
150        for (plan, s) in &by_plan {
151            let short = &plan[..plan.len().min(12)];
152            if s.runs >= min_runs
153                && (s.failed as f64 / s.runs as f64) >= min_failure_ratio
154            {
155                let mut args = Map::new();
156                args.insert("workflow".into(), json!(short));
157                args.insert("failed".into(), json!(s.failed));
158                args.insert("runs".into(), json!(s.runs));
159                args.insert(
160                    "rate".into(),
161                    json!(((s.failed as f64 / s.runs as f64) * 100.0).round() as i64),
162                );
163                args.insert("last_error".into(), json!(s.last_detail));
164                let mut data = Map::new();
165                data.insert("plan_hash".into(), json!(plan));
166                data.insert("failed".into(), json!(s.failed));
167                data.insert("runs".into(), json!(s.runs));
168                out.push(
169                    RecDraft::new(
170                        format!("host:workflow/{plan}"),
171                        ActionKind::Flag,
172                        Summary::new("run.failures", args),
173                        Proposal::Data { data },
174                    )
175                    .severity(Severity::High)
176                    .evidence(evidence.get(plan).cloned().unwrap_or_default()),
177                );
178            }
179            if min_usd > 0 && s.usd_micros >= min_usd && s.runs > 0 {
180                let mut args = Map::new();
181                args.insert("workflow".into(), json!(short));
182                args.insert("runs".into(), json!(s.runs));
183                args.insert(
184                    "usd".into(),
185                    json!(format!("{:.2}", s.usd_micros as f64 / 1e6)),
186                );
187                args.insert(
188                    "avg_usd".into(),
189                    json!(format!("{:.2}", s.usd_micros as f64 / 1e6 / s.runs as f64)),
190                );
191                let mut data = Map::new();
192                data.insert("plan_hash".into(), json!(plan));
193                data.insert("usd_micros".into(), json!(s.usd_micros));
194                out.push(
195                    RecDraft::new(
196                        format!("host:workflow-cost/{plan}"),
197                        ActionKind::Flag,
198                        Summary::new("run.cost", args),
199                        Proposal::Data { data },
200                    )
201                    .severity(Severity::Medium)
202                    .evidence(evidence.get(plan).cloned().unwrap_or_default()),
203                );
204            }
205            // Context pressure. A fold is the runtime keeping a long agent
206            // alive, so one is not a problem — a workflow that needs one on
207            // EVERY run is telling you its nodes no longer fit the window, and
208            // that is a plan-shape decision a person makes: split the node,
209            // bound its tool results, or accept the summaries. Advisory only;
210            // there is nothing here for an apply to do automatically.
211            if s.runs > 0 && s.folds > 0 && s.folds >= min_folds_per_run * s.runs {
212                let mut args = Map::new();
213                args.insert("workflow".into(), json!(short));
214                args.insert("runs".into(), json!(s.runs));
215                args.insert("folds".into(), json!(s.folds));
216                args.insert(
217                    "avg_folds".into(),
218                    json!(format!("{:.1}", s.folds as f64 / s.runs as f64)),
219                );
220                let mut data = Map::new();
221                data.insert("plan_hash".into(), json!(plan));
222                data.insert("folds".into(), json!(s.folds));
223                data.insert("runs".into(), json!(s.runs));
224                out.push(
225                    RecDraft::new(
226                        format!("host:workflow-context/{plan}"),
227                        ActionKind::Flag,
228                        Summary::new("run.context_pressure", args),
229                        Proposal::Data { data },
230                    )
231                    .severity(Severity::Medium)
232                    .evidence(evidence.get(plan).cloned().unwrap_or_default()),
233                );
234            }
235        }
236        Ok(out)
237    }
238
239    fn manifest(&self) -> &AnalyzerManifest {
240        &self.manifest
241    }
242}
243
244#[cfg(test)]
245mod tests {
246    use super::*;
247    use crate::testkit::TestSubstrate;
248
249    fn outcome(sub: &mut TestSubstrate, plan: &str, outcome: &str, usd: i64) {
250        sub.put_observation(
251            HARNESS_NS,
252            &[
253                ("observation_kind", json!("run_outcome")),
254                ("plan_hash", json!(plan)),
255                ("object", json!(outcome)),
256                ("spent_usd_micros", json!(usd)),
257                ("outcome_detail", json!("greet: ExecutorError: down")),
258            ],
259        );
260    }
261
262    /// A completed run that had to summarize itself `folds` times.
263    fn folded_run(sub: &mut TestSubstrate, plan: &str, folds: i64) {
264        sub.put_observation(
265            HARNESS_NS,
266            &[
267                ("observation_kind", json!("run_outcome")),
268                ("plan_hash", json!(plan)),
269                ("object", json!("completed")),
270                ("spent_usd_micros", json!(0)),
271                ("folds", json!(folds)),
272            ],
273        );
274    }
275
276    #[test]
277    fn flags_failing_workflow() {
278        let mut sub = TestSubstrate::new();
279        for _ in 0..2 {
280            outcome(&mut sub, "abcd1234", "failed", 0);
281        }
282        outcome(&mut sub, "abcd1234", "completed", 0);
283        let drafts = sub.analyze(&RunOutcome::new(), 10_000);
284        assert_eq!(drafts.len(), 1, "{drafts:?}");
285        let text = drafts[0].summary.render();
286        assert!(text.contains("failed 2/3"), "{text}");
287        assert_eq!(drafts[0].action_kind, ActionKind::Flag);
288    }
289
290    /// A workflow that folds on every run has outgrown the window. That is a
291    /// plan-shape decision (split the node, bound its tool results, or accept
292    /// the summaries) so it surfaces as an advisory Flag with nothing to
293    /// auto-apply — the same posture as the cost signal beside it.
294    #[test]
295    fn a_workflow_that_folds_every_run_is_surfaced() {
296        let mut sub = TestSubstrate::new();
297        for _ in 0..3 {
298            folded_run(&mut sub, "c0ffee11", 2);
299        }
300        let drafts = sub.analyze(&RunOutcome::new(), 10_000);
301        let pressure: Vec<_> = drafts
302            .iter()
303            .filter(|d| d.summary.render().contains("summarized its own transcript"))
304            .collect();
305        assert_eq!(pressure.len(), 1, "{drafts:?}");
306        let text = pressure[0].summary.render();
307        assert!(text.contains("6 time(s) across 3 runs"), "{text}");
308        assert!(text.contains("avg 2.0/run"), "{text}");
309        assert_eq!(pressure[0].action_kind, ActionKind::Flag);
310    }
311
312    /// Folding is the mechanism WORKING. An occasional fold must not nag: the
313    /// signal is a workflow that needs one every time, not one that ever needed
314    /// one at all.
315    #[test]
316    fn an_occasional_fold_is_not_a_finding() {
317        let mut sub = TestSubstrate::new();
318        folded_run(&mut sub, "c0ffee11", 1);
319        for _ in 0..5 {
320            outcome(&mut sub, "c0ffee11", "completed", 0);
321        }
322        let drafts = sub.analyze(&RunOutcome::new(), 10_000);
323        assert!(
324            !drafts.iter().any(|d| d.summary.render().contains("summarized its own")),
325            "one fold in six runs is healthy: {drafts:?}"
326        );
327    }
328
329    /// The field is absent on every run that never folded, and an absent field
330    /// must read as zero rather than as anything at all.
331    #[test]
332    fn runs_that_never_folded_carry_no_pressure_signal() {
333        let mut sub = TestSubstrate::new();
334        for _ in 0..4 {
335            outcome(&mut sub, "c0ffee11", "completed", 0);
336        }
337        let drafts = sub.analyze(&RunOutcome::new(), 10_000);
338        assert!(
339            !drafts.iter().any(|d| d.summary.render().contains("summarized its own")),
340            "{drafts:?}"
341        );
342    }
343
344    #[test]
345    fn healthy_workflow_stays_quiet() {
346        let mut sub = TestSubstrate::new();
347        for _ in 0..5 {
348            outcome(&mut sub, "abcd1234", "completed", 100);
349        }
350        assert!(sub.analyze(&RunOutcome::new(), 10_000).is_empty());
351    }
352
353    #[test]
354    fn cost_attribution_surfaces_expensive_workflows() {
355        let mut sub = TestSubstrate::new();
356        for _ in 0..4 {
357            outcome(&mut sub, "eeff5566", "completed", 2_000_000); // $2 each
358        }
359        let drafts = sub.analyze(&RunOutcome::new(), 10_000);
360        assert_eq!(drafts.len(), 1);
361        let text = drafts[0].summary.render();
362        assert!(text.contains("8.00"), "aggregate spend rendered: {text}");
363    }
364}