1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
//! The `Analyzer` trait and `AnalyzeCtx` — the SDK seam. `AnalyzeCtx` is a
//! struct (not a trait) so the engine can add methods without breaking
//! implementors. It exposes only read-only substrate access plus resolved
//! params, the watermark, and `now` — analyzers cannot write (trust floor),
//! enforced by holding `&dyn SubstrateRead`.
use crate::error::Result;
use crate::manifest::{AnalyzerManifest, Params};
use crate::model::GrainRecord;
use crate::recommendation::RecDraft;
use crate::substrate::{HeadGroup, ReadOpts, SubstrateRead, TelemetryView};
/// One applied recommendation due for outcome review, with its metric already
/// re-measured by the engine (which owns the `&mut` substrate). The outcome
/// analyzer makes the deterministic changed/regressed decision over this — I/O
/// in the engine, judgment in the analyzer.
#[derive(Debug, Clone, PartialEq)]
pub struct OutcomeInput {
pub rec_hash: String,
pub target_ref: String,
pub metric: String,
pub baseline: f64,
pub current: f64,
pub unit: String,
/// Carried from the metric snapshot so the analyzer applies the SAME
/// direction the engine did — see `recommendation::is_regression`.
pub higher_is_better: bool,
/// Where `baseline` came from (`OutcomeResult::baseline_kind`).
pub baseline_kind: String,
/// The evalset run `baseline` was read from, when it was read from one.
pub baseline_run_id: Option<String>,
/// The best value before the apply (`OutcomeResult::best_before`).
pub best_before: Option<f64>,
/// The minimum effect size the engine judged under, in the metric's unit
/// — passed through so the revert draft applies the SAME floor.
pub tolerance: f64,
/// The evalset run `current` was read from, when it was read from one.
pub current_run_id: Option<String>,
/// The cost bound's reading, when the policy set one
/// (`OutcomeResult::cost`). A breached bound with quality held is the
/// advisory-Flag case; with quality regressed it rides on the revert.
pub cost: Option<crate::recommendation::CostRead>,
}
/// The context handed to `analyze`. Read-only by construction.
pub struct AnalyzeCtx<'a> {
reader: &'a dyn SubstrateRead,
params: &'a Params,
namespaces: &'a [String],
watermark_ms: Option<i64>,
now_ms: i64,
outcome_inputs: &'a [OutcomeInput],
verdicts: &'a std::collections::BTreeMap<String, String>,
decider: Option<&'a crate::decide::Decider>,
}
impl<'a> AnalyzeCtx<'a> {
#[allow(clippy::too_many_arguments)]
pub fn new(
reader: &'a dyn SubstrateRead,
params: &'a Params,
namespaces: &'a [String],
watermark_ms: Option<i64>,
now_ms: i64,
outcome_inputs: &'a [OutcomeInput],
verdicts: &'a std::collections::BTreeMap<String, String>,
) -> Self {
AnalyzeCtx {
reader,
params,
namespaces,
watermark_ms,
now_ms,
outcome_inputs,
verdicts,
decider: None,
}
}
/// Hand the analyzer the engine's decision backend (or none). The engine
/// sets it on a production pass and never on a replay.
pub fn with_decider(mut self, decider: Option<&'a crate::decide::Decider>) -> Self {
self.decider = decider;
self
}
/// The installed decision backend, if any. An analyzer that uses one must
/// check [`crate::decide::Decider::calibrated`] before letting a
/// probability propose or drop anything, and must treat an `Err` as "no
/// contribution" (fail-soft), never as a failed analysis.
pub fn decider(&self) -> Option<&'a crate::decide::Decider> {
self.decider
}
/// The latest Verify-gate verdict for a grain an applied recommendation
/// created (`held`, `regressed`, `drifted`, `held_costlier`,
/// `unmeasured`), or `None` for a grain no apply created.
pub fn verdict_for(&self, created_hash: &str) -> Option<&str> {
self.verdicts.get(created_hash).map(String::as_str)
}
pub fn params(&self) -> &Params {
self.params
}
pub fn now_ms(&self) -> i64 {
self.now_ms
}
pub fn watermark_ms(&self) -> Option<i64> {
self.watermark_ms
}
pub fn capabilities(&self) -> crate::substrate::Capabilities {
self.reader.capabilities()
}
pub fn outcome_inputs(&self) -> &[OutcomeInput] {
self.outcome_inputs
}
/// All grains of a type across the configured namespaces (or all namespaces
/// when none are configured).
pub fn grains_of_type(&self, grain_type: &str, opts: ReadOpts) -> Result<Vec<GrainRecord>> {
if self.namespaces.is_empty() {
self.reader.grains_of_type(grain_type, None, opts)
} else {
let mut out = Vec::new();
for ns in self.namespaces {
out.extend(self.reader.grains_of_type(grain_type, Some(ns), opts)?);
}
Ok(out)
}
}
/// Live facts across configured namespaces.
pub fn facts(&self) -> Result<Vec<GrainRecord>> {
self.grains_of_type(crate::model::grain_type::FACT, ReadOpts::default())
}
/// Live observations across configured namespaces.
pub fn observations(&self) -> Result<Vec<GrainRecord>> {
self.grains_of_type(crate::model::grain_type::OBSERVATION, ReadOpts::default())
}
/// Live grains of one type in ONE explicit namespace — for analyzers
/// whose datasource lives in a well-known system namespace (the run
/// journals in `agent:harness`) regardless of the sweep's configured
/// scope.
pub fn grains_in(&self, grain_type: &str, namespace: &str) -> Result<Vec<GrainRecord>> {
self.reader
.grains_of_type(grain_type, Some(namespace), ReadOpts::default())
}
/// Live Skill grains across configured namespaces.
pub fn skills(&self) -> Result<Vec<GrainRecord>> {
self.grains_of_type(crate::model::grain_type::SKILL, ReadOpts::default())
}
/// Live Goal grains across configured namespaces.
pub fn goals(&self) -> Result<Vec<GrainRecord>> {
self.grains_of_type(crate::model::grain_type::GOAL, ReadOpts::default())
}
/// Tool grains (captured tool calls), optionally windowed by a `since`
/// watermark, live only. The flagship analyzer's input.
pub fn tools_since(&self, since_ms: Option<i64>) -> Result<Vec<GrainRecord>> {
self.grains_of_type(
crate::model::grain_type::TOOL,
ReadOpts {
live_only: true,
since_ms,
},
)
}
/// Entities with more than one live head (requires the forks capability).
pub fn heads(&self) -> Result<Vec<HeadGroup>> {
let ns = if self.namespaces.len() == 1 {
Some(self.namespaces[0].as_str())
} else {
None
};
self.reader.heads(ns)
}
/// A snapshot of the recall-telemetry rollups (requires the `telemetry`
/// capability). `None` when the substrate has no sidecar — a telemetry-fed
/// analyzer then degrades to an activation-ladder entry.
pub fn telemetry(&self) -> Result<Option<TelemetryView>> {
let ns = if self.namespaces.len() == 1 {
Some(self.namespaces[0].as_str())
} else {
None
};
self.reader.telemetry(ns)
}
}
/// An analysis unit. Object-safe so `builtin_analyzers()` yields trait objects.
pub trait Analyzer: Send + Sync {
fn manifest(&self) -> &AnalyzerManifest;
/// Produce recommendation drafts. `dedup_key`, `origin`, and the params
/// snapshot are stamped by the engine afterward, not here. Returning an
/// error drops *this* analyzer's findings for the run; other analyzers are
/// unaffected.
fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>>;
}
/// The default-registered built-in analyzers. Count is test-pinned (§11).
pub fn builtin_analyzers() -> Vec<Box<dyn Analyzer>> {
vec![
Box::new(crate::analyzers::tool_failure::ToolFailureClustering::new()),
Box::new(crate::analyzers::duplicate_sweep::DuplicateSweep::new()),
Box::new(crate::analyzers::contradiction_sweep::ContradictionSweep::new()),
Box::new(crate::analyzers::fork_surfacing::ForkSurfacing::new()),
Box::new(crate::analyzers::staleness::Staleness::new()),
Box::new(crate::analyzers::skill_stall::SkillStall::new()),
Box::new(crate::analyzers::goal_stagnation::GoalStagnation::new()),
Box::new(crate::analyzers::cold_grains::ColdGrains::new()),
Box::new(crate::analyzers::coverage_gap::CoverageGap::new()),
Box::new(crate::analyzers::budget_pressure::BudgetPressure::new()),
Box::new(crate::analyzers::outcome_review::OutcomeReview::new()),
Box::new(crate::analyzers::retention_sweep::RetentionSweep::new()),
Box::new(crate::analyzers::run_outcome::RunOutcome::new()),
Box::new(crate::analyzers::adapter_intake::AdapterIntake::new()),
Box::new(crate::analyzers::lesson_pile::LessonPile::new()),
]
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn builtins_have_unique_ids() {
let a = builtin_analyzers();
assert_eq!(
a.len(),
15,
"6 hygiene + skill/goal trajectory + 3 telemetry-fed (cold/coverage/budget) \
+ retention (default-off) + run_outcome + adapter_intake + lesson_pile (default-off)"
);
let mut ids: Vec<&str> = a.iter().map(|x| x.manifest().id.as_str()).collect();
ids.sort_unstable();
ids.dedup();
assert_eq!(ids.len(), 15, "analyzer ids must be unique");
}
#[test]
fn all_builtins_are_trust_class_builtin() {
for a in builtin_analyzers() {
assert_eq!(
a.manifest().trust_class,
crate::manifest::TrustClass::Builtin,
"{} must be builtin trust class",
a.manifest().id
);
}
}
}