1use 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
21const 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, 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 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); }
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}