Skip to main content

areev_loop/
analyzer.rs

1//! The `Analyzer` trait and `AnalyzeCtx` — the SDK seam. `AnalyzeCtx` is a
2//! struct (not a trait) so the engine can add methods without breaking
3//! implementors. It exposes only read-only substrate access plus resolved
4//! params, the watermark, and `now` — analyzers cannot write (trust floor),
5//! enforced by holding `&dyn SubstrateRead`.
6
7use 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/// One applied recommendation due for outcome review, with its metric already
14/// re-measured by the engine (which owns the `&mut` substrate). The outcome
15/// analyzer makes the deterministic changed/regressed decision over this — I/O
16/// in the engine, judgment in the analyzer.
17#[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    /// Carried from the metric snapshot so the analyzer applies the SAME
26    /// direction the engine did — see `recommendation::is_regression`.
27    pub higher_is_better: bool,
28    /// Where `baseline` came from (`OutcomeResult::baseline_kind`).
29    pub baseline_kind: String,
30    /// The evalset run `baseline` was read from, when it was read from one.
31    pub baseline_run_id: Option<String>,
32    /// The best value before the apply (`OutcomeResult::best_before`).
33    pub best_before: Option<f64>,
34    /// The minimum effect size the engine judged under, in the metric's unit
35    /// — passed through so the revert draft applies the SAME floor.
36    pub tolerance: f64,
37    /// The evalset run `current` was read from, when it was read from one.
38    pub current_run_id: Option<String>,
39    /// The cost bound's reading, when the policy set one
40    /// (`OutcomeResult::cost`). A breached bound with quality held is the
41    /// advisory-Flag case; with quality regressed it rides on the revert.
42    pub cost: Option<crate::recommendation::CostRead>,
43}
44
45/// The context handed to `analyze`. Read-only by construction.
46pub 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    /// Hand the analyzer the engine's decision backend (or none). The engine
81    /// sets it on a production pass and never on a replay.
82    pub fn with_decider(mut self, decider: Option<&'a crate::decide::Decider>) -> Self {
83        self.decider = decider;
84        self
85    }
86
87    /// The installed decision backend, if any. An analyzer that uses one must
88    /// check [`crate::decide::Decider::calibrated`] before letting a
89    /// probability propose or drop anything, and must treat an `Err` as "no
90    /// contribution" (fail-soft), never as a failed analysis.
91    pub fn decider(&self) -> Option<&'a crate::decide::Decider> {
92        self.decider
93    }
94
95    /// The latest Verify-gate verdict for a grain an applied recommendation
96    /// created (`held`, `regressed`, `drifted`, `held_costlier`,
97    /// `unmeasured`), or `None` for a grain no apply created.
98    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    /// All grains of a type across the configured namespaces (or all namespaces
119    /// when none are configured).
120    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    /// Live facts across configured namespaces.
133    pub fn facts(&self) -> Result<Vec<GrainRecord>> {
134        self.grains_of_type(crate::model::grain_type::FACT, ReadOpts::default())
135    }
136
137    /// Live observations across configured namespaces.
138    pub fn observations(&self) -> Result<Vec<GrainRecord>> {
139        self.grains_of_type(crate::model::grain_type::OBSERVATION, ReadOpts::default())
140    }
141
142    /// Live grains of one type in ONE explicit namespace — for analyzers
143    /// whose datasource lives in a well-known system namespace (the run
144    /// journals in `agent:harness`) regardless of the sweep's configured
145    /// scope.
146    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    /// Live Skill grains across configured namespaces.
152    pub fn skills(&self) -> Result<Vec<GrainRecord>> {
153        self.grains_of_type(crate::model::grain_type::SKILL, ReadOpts::default())
154    }
155
156    /// Live Goal grains across configured namespaces.
157    pub fn goals(&self) -> Result<Vec<GrainRecord>> {
158        self.grains_of_type(crate::model::grain_type::GOAL, ReadOpts::default())
159    }
160
161    /// Tool grains (captured tool calls), optionally windowed by a `since`
162    /// watermark, live only. The flagship analyzer's input.
163    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    /// Entities with more than one live head (requires the forks capability).
174    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    /// A snapshot of the recall-telemetry rollups (requires the `telemetry`
184    /// capability). `None` when the substrate has no sidecar — a telemetry-fed
185    /// analyzer then degrades to an activation-ladder entry.
186    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
196/// An analysis unit. Object-safe so `builtin_analyzers()` yields trait objects.
197pub trait Analyzer: Send + Sync {
198    fn manifest(&self) -> &AnalyzerManifest;
199
200    /// Produce recommendation drafts. `dedup_key`, `origin`, and the params
201    /// snapshot are stamped by the engine afterward, not here. Returning an
202    /// error drops *this* analyzer's findings for the run; other analyzers are
203    /// unaffected.
204    fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>>;
205}
206
207/// The default-registered built-in analyzers. Count is test-pinned (§11).
208pub 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}