Skip to main content

areev_loop/analyzers/
adapter_intake.rs

1//! Adapter intake (T0) — the tuning seam's propose leg. `areev tune`
2//! registers a host-trained adapter as an `mg:adapter` Fact in
3//! `agent:harness` (the registry tuple, its Rule E1 evalset pin, and the
4//! corpus-manifest lineage all embedded in the object JSON); this analyzer
5//! turns each unpromoted candidate into an `adapter_revision`
6//! recommendation. The engine's gates do the rest: apply is refused without
7//! a clean recorded run of the pinned evalset, the promotion is an
8//! immutable `(model:X, mg:adapter_promotion)` Fact hosts re-resolve from,
9//! and rollback retracts it.
10//!
11//! Lifecycle (deliberate, dedup-shaped): **one candidate per served model**.
12//! A subject with a live promotion is skipped entirely — replacing a
13//! promoted adapter means rolling the promotion back first, which frees the
14//! dedup key; the next pass then proposes the newest unpromoted candidate.
15//! A rolled-back candidate is re-proposed while its registry grain stays
16//! live ("the situation returned"); retiring the `mg:adapter` grain is how
17//! a host silences it.
18
19use 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
27/// The namespace `areev tune` registers adapters into.
28const 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            // No harness namespace / no read grant: nothing tuned — degrade
69            // to nothing, never fabricate.
70            Err(_) => return Ok(Vec::new()),
71        };
72        // Subjects with a live promotion are settled: one candidate per
73        // served model, and replacement starts with a rollback (which
74        // retracts the promotion Fact, emptying this set for the subject).
75        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        // Newest live candidate per subject (created_at, hash tiebreak) —
85        // the engine's in-pass dedup would otherwise keep whichever came
86        // first, which after a rollback is the OLD adapter.
87        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            // The registry tuple: everything the promotion needs rides in the
110            // object JSON `areev tune` validated at record time. A grain that
111            // does not parse, or that carries no pin, is not a candidate —
112            // Rule E1 would refuse the draft at the door anyway; skipping it
113            // here keeps the queue free of dead-on-arrival entries.
114            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            // The verify leg: after promotion the host re-runs the pinned
146            // evalset against the served adapter; a recorded run with more
147            // failures than the gating baseline (0 — apply refuses a failing
148            // gate) is a regression and outcome_review proposes the revert.
149            // No recorded baseline run yet → no metric: honestly unmeasured,
150            // never fabricated.
151            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                    // Apply refuses a failing gate, so the promoted state's
176                    // baseline is zero failures; any post-apply failure is a
177                    // regression.
178                    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                    // 1 day, 1 week, 1 month — a late regression is caught by
190                    // the schedule.
191                    horizons_ms: vec![86_400_000, 7 * 86_400_000, 30 * 86_400_000],
192                    checkpoints: Vec::new(),
193                    // Failure count: fewer is better (the default).
194                    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        // The third lock: the target class itself is auto-apply-ineligible.
235        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        // No recorded eval run for the pin yet: honestly unmeasured.
253        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        // One candidate per served model: a live promotion settles the slot.
285        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}