use crate::analyzer::{AnalyzeCtx, Analyzer};
use crate::error::Result;
use crate::manifest::*;
use crate::model::{ActionKind, GrainRecord, Severity};
use crate::recommendation::{MetricSnapshot, Proposal, RecDraft, Summary};
use serde_json::{json, Map, Value};
use std::collections::{BTreeMap, BTreeSet};
const HARNESS_NS: &str = "agent:harness";
pub struct AdapterIntake {
manifest: AnalyzerManifest,
}
impl AdapterIntake {
pub fn new() -> Self {
AdapterIntake {
manifest: AnalyzerManifest {
id: "loop.adapter_intake/1".into(),
title: "Adapter intake".into(),
description:
"Proposes promoting host-trained adapters registered by `areev tune` \
— one candidate per served model, gated on the pinned evalset, \
never auto-applied."
.into(),
tier: Tier::T0,
cadence: CadenceClass::Slow,
requires: vec![],
target_classes: vec![TargetClass::Host],
auto_apply: AutoApplyClass::Never,
trust_class: TrustClass::Builtin,
params: vec![],
default_on: true,
},
}
}
}
impl Default for AdapterIntake {
fn default() -> Self {
Self::new()
}
}
impl Analyzer for AdapterIntake {
fn analyze(&self, ctx: &AnalyzeCtx) -> Result<Vec<RecDraft>> {
let harness = match ctx.grains_in("fact", HARNESS_NS) {
Ok(rows) => rows,
Err(_) => return Ok(Vec::new()),
};
let promoted: BTreeSet<String> = match ctx.grains_in("fact", crate::engine::LOOP_NS) {
Ok(rows) => rows
.iter()
.filter(|g| g.is_live() && g.str_field("relation") == Some("mg:adapter_promotion"))
.filter_map(|g| g.str_field("subject").map(str::to_string))
.collect(),
Err(_) => BTreeSet::new(),
};
let mut newest: BTreeMap<String, &GrainRecord> = BTreeMap::new();
for g in &harness {
if !g.is_live() || g.str_field("relation") != Some("mg:adapter") {
continue;
}
let Some(subject) = g.str_field("subject") else {
continue;
};
if promoted.contains(subject) {
continue;
}
let newer = match newest.get(subject) {
Some(cur) => (g.created_at_ms, g.hash.as_str()) > (cur.created_at_ms, cur.hash.as_str()),
None => true,
};
if newer {
newest.insert(subject.to_string(), g);
}
}
let mut drafts = Vec::new();
for (subject, grain) in newest {
let Some(tuple) = grain
.str_field("object")
.and_then(|o| serde_json::from_str::<Map<String, Value>>(o).ok())
else {
continue;
};
let Some(pin) = tuple
.get("evalset_hash")
.and_then(Value::as_str)
.filter(|h| !h.trim().is_empty())
.map(str::to_string)
else {
continue;
};
let base_model = tuple
.get("base_model")
.and_then(Value::as_str)
.unwrap_or("?")
.to_string();
let mut evidence = vec![grain.hash.clone()];
if let Some(manifest) = tuple.get("corpus_manifest").and_then(Value::as_str) {
evidence.push(manifest.to_string());
}
let mut args = Map::new();
args.insert("model".into(), json!(subject));
args.insert("base_model".into(), json!(base_model));
let mut payload = tuple.clone();
payload.insert("adapter_grain".into(), json!(grain.hash));
let eval_subject = format!("evalset:{pin}");
let baseline_run = harness
.iter()
.filter(|g| {
g.str_field("relation") == Some("mg:eval_run")
&& g.str_field("subject") == Some(eval_subject.as_str())
})
.max_by(|a, b| {
(a.created_at_ms, a.hash.as_str()).cmp(&(b.created_at_ms, b.hash.as_str()))
});
let mut draft = RecDraft::new(
subject.clone(),
ActionKind::AdapterRevision,
Summary::new("adapter.candidate", args),
Proposal::Data { data: payload },
)
.severity(Severity::Medium)
.evidence(evidence)
.evalset_hash(pin.clone());
if baseline_run.is_some() {
draft = draft.metric(MetricSnapshot {
metric: format!("evalset:{pin}:failed"),
baseline: 0.0,
unit: "count".into(),
n: 1,
window: "per-run".into(),
subject: Some(subject.clone()),
namespace: None,
relation: None,
query: format!(
"RECALL facts WHERE subject = \"evalset:{pin}\" AND relation = \"mg:eval_run\""
),
review_after_ms: 86_400_000,
horizons_ms: vec![86_400_000, 7 * 86_400_000, 30 * 86_400_000],
checkpoints: Vec::new(),
higher_is_better: false,
});
}
drafts.push(draft);
}
Ok(drafts)
}
fn manifest(&self) -> &AnalyzerManifest {
&self.manifest
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::model::TargetRef;
use crate::testkit::TestSubstrate;
const T: i64 = 1_000_000;
fn tuple(pin: &str, manifest: &str) -> String {
json!({
"adapter": {"uri": "file:///adapters/a.safetensors", "sha256": "feed"},
"base_model": "qwen3-4b",
"quantization": "bf16",
"serving_runtime": "vllm",
"serves_as": "acme-support",
"evalset_hash": pin,
"corpus_manifest": manifest,
})
.to_string()
}
#[test]
fn manifest_is_never_auto_apply_and_the_class_is_ineligible() {
let a = AdapterIntake::new();
assert_eq!(a.manifest().id, "loop.adapter_intake/1");
assert_eq!(a.manifest().auto_apply, AutoApplyClass::Never);
assert!(a.manifest().default_on);
assert!(!TargetRef::parse("model:acme-support")
.unwrap()
.auto_apply_eligible_class());
}
#[test]
fn drafts_the_candidate_with_pin_and_evidence() {
let mut sub = TestSubstrate::new();
let adapter =
sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
let drafts = sub.analyze(&AdapterIntake::new(), T + 1);
assert_eq!(drafts.len(), 1);
let d = &drafts[0];
assert_eq!(d.target_ref, "model:acme-support");
assert_eq!(d.action_kind, ActionKind::AdapterRevision);
assert_eq!(d.evalset_hash.as_deref(), Some("pin1"));
assert_eq!(d.evidence, vec![adapter, "cafe".to_string()]);
assert!(d.metric.is_none());
match &d.proposal {
Proposal::Data { data } => {
assert_eq!(data.get("serves_as").and_then(Value::as_str), Some("acme-support"));
assert!(data.get("adapter_grain").is_some());
}
other => panic!("expected Data proposal, got {other:?}"),
}
}
#[test]
fn attaches_the_evalset_metric_when_a_baseline_run_exists() {
let mut sub = TestSubstrate::new();
sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
sub.add_fact_at(
"agent:harness",
"evalset:pin1",
"mg:eval_run",
"{\"run_id\":\"eval-1\",\"passed\":3,\"failed\":0}",
T,
);
let drafts = sub.analyze(&AdapterIntake::new(), T + 1);
let m = drafts[0].metric.as_ref().expect("the verify metric");
assert_eq!(m.metric, "evalset:pin1:failed");
assert!(!m.higher_is_better, "failure counts: fewer is better");
}
#[test]
fn a_promoted_subject_is_silent_until_rollback_retracts_the_promotion() {
let mut sub = TestSubstrate::new();
sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "cafe"), T);
sub.add_fact_at("areev-loop", "model:acme-support", "mg:adapter_promotion", "{}", T + 10);
assert!(sub.analyze(&AdapterIntake::new(), T + 20).is_empty());
}
#[test]
fn two_candidates_one_model_yields_only_the_newest() {
let mut sub = TestSubstrate::new();
sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "old"), T);
let newer =
sub.add_fact_at("agent:harness", "model:acme-support", "mg:adapter", &tuple("pin1", "new"), T + 100);
let drafts = sub.analyze(&AdapterIntake::new(), T + 200);
assert_eq!(drafts.len(), 1);
assert_eq!(drafts[0].evidence[0], newer);
}
#[test]
fn malformed_or_unpinned_registry_grains_are_skipped() {
let mut sub = TestSubstrate::new();
sub.add_fact_at("agent:harness", "model:a", "mg:adapter", "not json", T);
sub.add_fact_at(
"agent:harness",
"model:b",
"mg:adapter",
"{\"base_model\":\"x\",\"serves_as\":\"b\"}",
T,
);
assert!(sub.analyze(&AdapterIntake::new(), T + 1).is_empty());
}
}