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 checkpoints: Vec::new(),
193 higher_is_better: false,
195 });
196 }
197 drafts.push(draft);
198 }
199 Ok(drafts)
200 }
201
202 fn manifest(&self) -> &AnalyzerManifest {
203 &self.manifest
204 }
205}
206
207#[cfg(test)]
208mod tests {
209 use super::*;
210 use crate::model::TargetRef;
211 use crate::testkit::TestSubstrate;
212
213 const T: i64 = 1_000_000;
214
215 fn tuple(pin: &str, manifest: &str) -> String {
216 json!({
217 "adapter": {"uri": "file:///adapters/a.safetensors", "sha256": "feed"},
218 "base_model": "qwen3-4b",
219 "quantization": "bf16",
220 "serving_runtime": "vllm",
221 "serves_as": "acme-support",
222 "evalset_hash": pin,
223 "corpus_manifest": manifest,
224 })
225 .to_string()
226 }
227
228 #[test]
229 fn manifest_is_never_auto_apply_and_the_class_is_ineligible() {
230 let a = AdapterIntake::new();
231 assert_eq!(a.manifest().id, "loop.adapter_intake/1");
232 assert_eq!(a.manifest().auto_apply, AutoApplyClass::Never);
233 assert!(a.manifest().default_on);
234 assert!(!TargetRef::parse("model:acme-support")
236 .unwrap()
237 .auto_apply_eligible_class());
238 }
239
240 #[test]
241 fn drafts_the_candidate_with_pin_and_evidence() {
242 let mut sub = TestSubstrate::new();
243 let adapter =
244 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
245 let drafts = sub.analyze(&AdapterIntake::new(), T + 1);
246 assert_eq!(drafts.len(), 1);
247 let d = &drafts[0];
248 assert_eq!(d.target_ref, "model:acme-support");
249 assert_eq!(d.action_kind, ActionKind::AdapterRevision);
250 assert_eq!(d.evalset_hash.as_deref(), Some("pin1"));
251 assert_eq!(d.evidence, vec![adapter, "cafe".to_string()]);
252 assert!(d.metric.is_none());
254 match &d.proposal {
255 Proposal::Data { data } => {
256 assert_eq!(data.get("serves_as").and_then(Value::as_str), Some("acme-support"));
257 assert!(data.get("adapter_grain").is_some());
258 }
259 other => panic!("expected Data proposal, got {other:?}"),
260 }
261 }
262
263 #[test]
264 fn attaches_the_evalset_metric_when_a_baseline_run_exists() {
265 let mut sub = TestSubstrate::new();
266 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
267 sub.add_fact_at(
268 "agent:harness",
269 "evalset:pin1",
270 "mg:eval_run",
271 "{\"run_id\":\"eval-1\",\"passed\":3,\"failed\":0}",
272 T,
273 );
274 let drafts = sub.analyze(&AdapterIntake::new(), T + 1);
275 let m = drafts[0].metric.as_ref().expect("the verify metric");
276 assert_eq!(m.metric, "evalset:pin1:failed");
277 assert!(!m.higher_is_better, "failure counts: fewer is better");
278 }
279
280 #[test]
281 fn a_promoted_subject_is_silent_until_rollback_retracts_the_promotion() {
282 let mut sub = TestSubstrate::new();
283 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
284 sub.add_fact_at("areev-loop", "model:acme-support", "mg:adapter_promotion", "{}", T + 10);
286 assert!(sub.analyze(&AdapterIntake::new(), T + 20).is_empty());
287 }
288
289 #[test]
290 fn two_candidates_one_model_yields_only_the_newest() {
291 let mut sub = TestSubstrate::new();
292 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "old"), T);
293 let newer =
294 sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "new"), T + 100);
295 let drafts = sub.analyze(&AdapterIntake::new(), T + 200);
296 assert_eq!(drafts.len(), 1);
297 assert_eq!(drafts[0].evidence[0], newer);
298 }
299
300 #[test]
301 fn malformed_or_unpinned_registry_grains_are_skipped() {
302 let mut sub = TestSubstrate::new();
303 sub.add_fact_at("agent:harness", "model:a", "mg:adapter", "not json", T);
304 sub.add_fact_at(
305 "agent:harness",
306 "model:b",
307 "mg:adapter",
308 "{\"base_model\":\"x\",\"serves_as\":\"b\"}",
309 T,
310 );
311 assert!(sub.analyze(&AdapterIntake::new(), T + 1).is_empty());
312 }
313}