1use crate::analyzer::{AnalyzeCtx, Analyzer};
20use crate::error::Result;
21use crate::manifest::*;
22use crate::model::{ActionKind, GrainRecord, Severity};
23use crate::recommendation::{MetricSnapshot, Proposal, RecDraft, Summary};
24use serde_json::{json, Map, Value};
25use std::collections::{BTreeMap, BTreeSet};
26
27const HARNESS_NS: &str = "agent:harness";
29
30pub struct AdapterIntake {
31 manifest: AnalyzerManifest,
32}
33
34impl AdapterIntake {
35 pub fn new() -> Self {
36 AdapterIntake {
37 manifest: AnalyzerManifest {
38 id: "loop.adapter_intake/1".into(),
39 title: "Adapter intake".into(),
40 description:
41 "Proposes promoting host-trained adapters registered by `areev tune` \
42 — one candidate per served model, gated on the pinned evalset, \
43 never auto-applied."
44 .into(),
45 tier: Tier::T0,
46 cadence: CadenceClass::Slow,
47 requires: vec![],
48 target_classes: vec![TargetClass::Host],
49 auto_apply: AutoApplyClass::Never,
50 trust_class: TrustClass::Builtin,
51 params: vec![],
52 default_on: true,
53 },
54 }
55 }
56}
57
58impl Default for AdapterIntake {
59 fn default() -> Self {
60 Self::new()
61 }
62}
63
64impl Analyzer for AdapterIntake {
65 fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
66 let harness = match ctx.grains_in("fact", HARNESS_NS) {
67 Ok(rows) => rows,
68 Err(_) => return Ok(Vec::new()),
71 };
72 let promoted: BTreeSet<String> = match ctx.grains_in("fact", crate::engine::LOOP_NS) {
76 Ok(rows) => rows
77 .iter()
78 .filter(|g| g.is_live() && g.str_field("relation") == Some("mg:adapter_promotion"))
79 .filter_map(|g| g.str_field("subject").map(str::to_string))
80 .collect(),
81 Err(_) => BTreeSet::new(),
82 };
83
84 let mut newest: BTreeMap<String, &GrainRecord> = BTreeMap::new();
88 for g in &harness {
89 if !g.is_live() || g.str_field("relation") != Some("mg:adapter") {
90 continue;
91 }
92 let Some(subject) = g.str_field("subject") else {
93 continue;
94 };
95 if promoted.contains(subject) {
96 continue;
97 }
98 let newer = match newest.get(subject) {
99 Some(cur) => (g.created_at_ms, g.hash.as_str()) > (cur.created_at_ms, cur.hash.as_str()),
100 None => true,
101 };
102 if newer {
103 newest.insert(subject.to_string(), g);
104 }
105 }
106
107 let mut drafts = Vec::new();
108 for (subject, grain) in newest {
109 let Some(tuple) = grain
115 .str_field("object")
116 .and_then(|o| serde_json::from_str::<Map<String, Value>>(o).ok())
117 else {
118 continue;
119 };
120 let Some(pin) = tuple
121 .get("evalset_hash")
122 .and_then(Value::as_str)
123 .filter(|h| !h.trim().is_empty())
124 .map(str::to_string)
125 else {
126 continue;
127 };
128 let base_model = tuple
129 .get("base_model")
130 .and_then(Value::as_str)
131 .unwrap_or("?")
132 .to_string();
133 let mut evidence = vec![grain.hash.clone()];
134 if let Some(manifest) = tuple.get("corpus_manifest").and_then(Value::as_str) {
135 evidence.push(manifest.to_string());
136 }
137
138 let mut args = Map::new();
139 args.insert("model".into(), json!(subject));
140 args.insert("base_model".into(), json!(base_model));
141
142 let mut payload = tuple.clone();
143 payload.insert("adapter_grain".into(), json!(grain.hash));
144
145 let eval_subject = format!("evalset:{pin}");
152 let baseline_run = harness
153 .iter()
154 .filter(|g| {
155 g.str_field("relation") == Some("mg:eval_run")
156 && g.str_field("subject") == Some(eval_subject.as_str())
157 })
158 .max_by(|a, b| {
159 (a.created_at_ms, a.hash.as_str()).cmp(&(b.created_at_ms, b.hash.as_str()))
160 });
161
162 let mut draft = RecDraft::new(
163 subject.clone(),
164 ActionKind::AdapterRevision,
165 Summary::new("adapter.candidate", args),
166 Proposal::Data { data: payload },
167 )
168 .severity(Severity::Medium)
169 .evidence(evidence)
170 .evalset_hash(pin.clone());
171
172 if baseline_run.is_some() {
173 draft = draft.metric(MetricSnapshot {
174 metric: format!("evalset:{pin}:failed"),
175 baseline: 0.0,
179 unit: "count".into(),
180 n: 1,
181 window: "per-run".into(),
182 subject: Some(subject.clone()),
183 namespace: None,
184 relation: None,
185 query: format!(
186 "RECALL facts WHERE subject = \"evalset:{pin}\" AND relation = \"mg:eval_run\""
187 ),
188 review_after_ms: 86_400_000,
189 horizons_ms: vec![86_400_000, 7 * 86_400_000, 30 * 86_400_000],
192 higher_is_better: false,
194 });
195 }
196 drafts.push(draft);
197 }
198 Ok(drafts)
199 }
200
201 fn manifest(&self) -> &AnalyzerManifest {
202 &self.manifest
203 }
204}
205
206#[cfg(test)]
207mod tests {
208 use super::*;
209 use crate::model::TargetRef;
210 use crate::testkit::TestSubstrate;
211
212 const T: i64 = 1_000_000;
213
214 fn tuple(pin: &str, manifest: &str) -> String {
215 json!({
216 "adapter": {"uri": "file:///adapters/a.safetensors", "sha256": "feed"},
217 "base_model": "qwen3-4b",
218 "quantization": "bf16",
219 "serving_runtime": "vllm",
220 "serves_as": "acme-support",
221 "evalset_hash": pin,
222 "corpus_manifest": manifest,
223 })
224 .to_string()
225 }
226
227 #[test]
228 fn manifest_is_never_auto_apply_and_the_class_is_ineligible() {
229 let a = AdapterIntake::new();
230 assert_eq!(a.manifest().id, "loop.adapter_intake/1");
231 assert_eq!(a.manifest().auto_apply, AutoApplyClass::Never);
232 assert!(a.manifest().default_on);
233 assert!(!TargetRef::parse("model:acme-support")
235 .unwrap()
236 .auto_apply_eligible_class());
237 }
238
239 #[test]
240 fn drafts_the_candidate_with_pin_and_evidence() {
241 let mut sub = TestSubstrate::new();
242 let adapter =
243 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
244 let drafts = sub.analyze(&AdapterIntake::new(), T + 1);
245 assert_eq!(drafts.len(), 1);
246 let d = &drafts[0];
247 assert_eq!(d.target_ref, "model:acme-support");
248 assert_eq!(d.action_kind, ActionKind::AdapterRevision);
249 assert_eq!(d.evalset_hash.as_deref(), Some("pin1"));
250 assert_eq!(d.evidence, vec![adapter, "cafe".to_string()]);
251 assert!(d.metric.is_none());
253 match &d.proposal {
254 Proposal::Data { data } => {
255 assert_eq!(data.get("serves_as").and_then(Value::as_str), Some("acme-support"));
256 assert!(data.get("adapter_grain").is_some());
257 }
258 other => panic!("expected Data proposal, got {other:?}"),
259 }
260 }
261
262 #[test]
263 fn attaches_the_evalset_metric_when_a_baseline_run_exists() {
264 let mut sub = TestSubstrate::new();
265 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
266 sub.add_fact_at(
267 "agent:harness",
268 "evalset:pin1",
269 "mg:eval_run",
270 "{\"run_id\":\"eval-1\",\"passed\":3,\"failed\":0}",
271 T,
272 );
273 let drafts = sub.analyze(&AdapterIntake::new(), T + 1);
274 let m = drafts[0].metric.as_ref().expect("the verify metric");
275 assert_eq!(m.metric, "evalset:pin1:failed");
276 assert!(!m.higher_is_better, "failure counts: fewer is better");
277 }
278
279 #[test]
280 fn a_promoted_subject_is_silent_until_rollback_retracts_the_promotion() {
281 let mut sub = TestSubstrate::new();
282 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
283 sub.add_fact_at("areev-loop", "model:acme-support", "mg:adapter_promotion", "{}", T + 10);
285 assert!(sub.analyze(&AdapterIntake::new(), T + 20).is_empty());
286 }
287
288 #[test]
289 fn two_candidates_one_model_yields_only_the_newest() {
290 let mut sub = TestSubstrate::new();
291 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "old"), T);
292 let newer =
293 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "new"), T + 100);
294 let drafts = sub.analyze(&AdapterIntake::new(), T + 200);
295 assert_eq!(drafts.len(), 1);
296 assert_eq!(drafts[0].evidence[0], newer);
297 }
298
299 #[test]
300 fn malformed_or_unpinned_registry_grains_are_skipped() {
301 let mut sub = TestSubstrate::new();
302 sub.add_fact_at("agent:harness", "model:a", "mg:adapter", "not json", T);
303 sub.add_fact_at(
304 "agent:harness",
305 "model:b",
306 "mg:adapter",
307 "{\"base_model\":\"x\",\"serves_as\":\"b\"}",
308 T,
309 );
310 assert!(sub.analyze(&AdapterIntake::new(), T + 1).is_empty());
311 }
312}