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, 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 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, 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 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 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 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 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 #[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 #[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 #[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); }
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}