1use crate::error::Result;
8use crate::manifest::{AnalyzerManifest, Params};
9use crate::model::GrainRecord;
10use crate::recommendation::RecDraft;
11use crate::substrate::{HeadGroup, ReadOpts, SubstrateRead, TelemetryView};
12
13#[derive(Debug, Clone, PartialEq)]
18pub struct OutcomeInput {
19 pub rec_hash: String,
20 pub target_ref: String,
21 pub metric: String,
22 pub baseline: f64,
23 pub current: f64,
24 pub unit: String,
25 pub higher_is_better: bool,
28 pub baseline_kind: String,
30 pub baseline_run_id: Option<String>,
32 pub best_before: Option<f64>,
34 pub tolerance: f64,
37 pub current_run_id: Option<String>,
39 pub cost: Option<crate::recommendation::CostRead>,
43}
44
45pub struct AnalyzeCtx<'a> {
47 reader: &'a dyn SubstrateRead,
48 params: &'a Params,
49 namespaces: &'a [String],
50 watermark_ms: Option<i64>,
51 now_ms: i64,
52 outcome_inputs: &'a [OutcomeInput],
53 verdicts: &'a std::collections::BTreeMap<String, String>,
54 decider: Option<&'a crate::decide::Decider>,
55}
56
57impl<'a> AnalyzeCtx<'a> {
58 #[allow(clippy::too_many_arguments)]
59 pub fn new(
60 reader: &'a dyn SubstrateRead,
61 params: &'a Params,
62 namespaces: &'a [String],
63 watermark_ms: Option<i64>,
64 now_ms: i64,
65 outcome_inputs: &'a [OutcomeInput],
66 verdicts: &'a std::collections::BTreeMap<String, String>,
67 ) -> Self {
68 AnalyzeCtx {
69 reader,
70 params,
71 namespaces,
72 watermark_ms,
73 now_ms,
74 outcome_inputs,
75 verdicts,
76 decider: None,
77 }
78 }
79
80 pub fn with_decider(mut self, decider: Option<&'a crate::decide::Decider>) -> Self {
83 self.decider = decider;
84 self
85 }
86
87 pub fn decider(&self) -> Option<&'a crate::decide::Decider> {
92 self.decider
93 }
94
95 pub fn verdict_for(&self, created_hash: &str) -> Option<&str> {
99 self.verdicts.get(created_hash).map(String::as_str)
100 }
101
102 pub fn params(&self) -> &Params {
103 self.params
104 }
105 pub fn now_ms(&self) -> i64 {
106 self.now_ms
107 }
108 pub fn watermark_ms(&self) -> Option<i64> {
109 self.watermark_ms
110 }
111 pub fn capabilities(&self) -> crate::substrate::Capabilities {
112 self.reader.capabilities()
113 }
114 pub fn outcome_inputs(&self) -> &[OutcomeInput] {
115 self.outcome_inputs
116 }
117
118 pub fn grains_of_type(&self, grain_type: &str, opts: ReadOpts) -> Result<Vec<GrainRecord>> {
121 if self.namespaces.is_empty() {
122 self.reader.grains_of_type(grain_type, None, opts)
123 } else {
124 let mut out = Vec::new();
125 for ns in self.namespaces {
126 out.extend(self.reader.grains_of_type(grain_type, Some(ns), opts)?);
127 }
128 Ok(out)
129 }
130 }
131
132 pub fn facts(&self) -> Result<Vec<GrainRecord>> {
134 self.grains_of_type(crate::model::grain_type::FACT, ReadOpts::default())
135 }
136
137 pub fn observations(&self) -> Result<Vec<GrainRecord>> {
139 self.grains_of_type(crate::model::grain_type::OBSERVATION, ReadOpts::default())
140 }
141
142 pub fn grains_in(&self, grain_type: &str, namespace: &str) -> Result<Vec<GrainRecord>> {
147 self.reader
148 .grains_of_type(grain_type, Some(namespace), ReadOpts::default())
149 }
150
151 pub fn skills(&self) -> Result<Vec<GrainRecord>> {
153 self.grains_of_type(crate::model::grain_type::SKILL, ReadOpts::default())
154 }
155
156 pub fn goals(&self) -> Result<Vec<GrainRecord>> {
158 self.grains_of_type(crate::model::grain_type::GOAL, ReadOpts::default())
159 }
160
161 pub fn tools_since(&self, since_ms: Option<i64>) -> Result<Vec<GrainRecord>> {
164 self.grains_of_type(
165 crate::model::grain_type::TOOL,
166 ReadOpts {
167 live_only: true,
168 since_ms,
169 },
170 )
171 }
172
173 pub fn heads(&self) -> Result<Vec<HeadGroup>> {
175 let ns = if self.namespaces.len() == 1 {
176 Some(self.namespaces[0].as_str())
177 } else {
178 None
179 };
180 self.reader.heads(ns)
181 }
182
183 pub fn telemetry(&self) -> Result<Option<TelemetryView>> {
187 let ns = if self.namespaces.len() == 1 {
188 Some(self.namespaces[0].as_str())
189 } else {
190 None
191 };
192 self.reader.telemetry(ns)
193 }
194}
195
196pub trait Analyzer: Send + Sync {
198 fn manifest(&self) -> &AnalyzerManifest;
199
200 fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>>;
205}
206
207pub fn builtin_analyzers() -> Vec<Box<dyn Analyzer>> {
209 vec![
210 Box::new(crate::analyzers::tool_failure::ToolFailureClustering::new()),
211 Box::new(crate::analyzers::duplicate_sweep::DuplicateSweep::new()),
212 Box::new(crate::analyzers::contradiction_sweep::ContradictionSweep::new()),
213 Box::new(crate::analyzers::fork_surfacing::ForkSurfacing::new()),
214 Box::new(crate::analyzers::staleness::Staleness::new()),
215 Box::new(crate::analyzers::skill_stall::SkillStall::new()),
216 Box::new(crate::analyzers::goal_stagnation::GoalStagnation::new()),
217 Box::new(crate::analyzers::cold_grains::ColdGrains::new()),
218 Box::new(crate::analyzers::coverage_gap::CoverageGap::new()),
219 Box::new(crate::analyzers::budget_pressure::BudgetPressure::new()),
220 Box::new(crate::analyzers::outcome_review::OutcomeReview::new()),
221 Box::new(crate::analyzers::retention_sweep::RetentionSweep::new()),
222 Box::new(crate::analyzers::run_outcome::RunOutcome::new()),
223 Box::new(crate::analyzers::adapter_intake::AdapterIntake::new()),
224 Box::new(crate::analyzers::lesson_pile::LessonPile::new()),
225 ]
226}
227
228#[cfg(test)]
229mod tests {
230 use super::*;
231
232 #[test]
233 fn builtins_have_unique_ids() {
234 let a = builtin_analyzers();
235 assert_eq!(
236 a.len(),
237 15,
238 "6 hygiene + skill/goal trajectory + 3 telemetry-fed (cold/coverage/budget) \
239 + retention (default-off) + run_outcome + adapter_intake + lesson_pile (default-off)"
240 );
241 let mut ids: Vec<&str> = a.iter().map(|x| x.manifest().id.as_str()).collect();
242 ids.sort_unstable();
243 ids.dedup();
244 assert_eq!(ids.len(), 15, "analyzer ids must be unique");
245 }
246
247 #[test]
248 fn all_builtins_are_trust_class_builtin() {
249 for a in builtin_analyzers() {
250 assert_eq!(
251 a.manifest().trust_class,
252 crate::manifest::TrustClass::Builtin,
253 "{} must be builtin trust class",
254 a.manifest().id
255 );
256 }
257 }
258}