use crate::{probe, setup, GlobalOpts};
use anyhow::{Context, Result};
use mecha_core::config::{Config, ProviderConfig};
use mecha_core::counterfactual::ProbeVerdict;
use mecha_core::eval::Judge;
use mecha_core::learning::{
domain_rules_section, locate_followup, rules_hash, strip_rules_block, wrap_rules_block,
LearningStore, Rule, Trigger, ValidationRecord,
};
use mecha_core::message::{CompletionRequest, Message};
use mecha_core::session::Session;
use std::collections::BTreeMap;
#[derive(clap::Args, Debug)]
pub struct Args {
#[arg(long)]
pub judge_provider: Option<String>,
#[arg(long)]
pub judge_model: Option<String>,
#[arg(long)]
pub unprocessed_only: bool,
#[arg(long, value_delimiter = ',')]
pub trigger: Vec<String>,
#[arg(long)]
pub no_attribute: bool,
}
struct RuleSurface {
flat: Vec<(String, Rule)>,
user_by_domain: BTreeMap<String, Vec<Rule>>,
}
impl RuleSurface {
fn load(store: &LearningStore) -> Result<Self> {
let mut flat = Vec::new();
let mut user_by_domain = BTreeMap::new();
for domain in store.domains() {
user_by_domain.insert(domain.clone(), store.user_rules(&domain)?);
for rule in store.learned_rules(&domain)? {
if rule.active() {
flat.push((domain.clone(), rule));
}
}
}
Ok(RuleSurface {
flat,
user_by_domain,
})
}
fn rule_ids(&self) -> Vec<String> {
self.flat.iter().filter_map(|(_, r)| r.id.clone()).collect()
}
fn block_with(&self, selected: &[usize]) -> Option<String> {
let mut sections = Vec::new();
for (domain, user) in &self.user_by_domain {
let learned: Vec<Rule> = self
.flat
.iter()
.enumerate()
.filter(|(i, (d, _))| d == domain && selected.contains(i))
.map(|(_, (_, r))| r.clone())
.collect();
sections.extend(domain_rules_section(domain, user, &learned));
}
wrap_rules_block(sections)
}
}
async fn attribute_regression(
prepared: &setup::Prepared,
provider_cfg: &ProviderConfig,
model: &str,
prep: &probe::ProbePrep,
surface: &RuleSurface,
) -> Result<Option<usize>> {
if surface.flat.is_empty() {
return Ok(None);
}
let fails = |selected: Vec<usize>| async move {
let block = surface.block_with(&selected);
match probe::drive_arm(prepared, provider_cfg, model, prep, block.as_deref()).await? {
Ok(ProbeVerdict::Fail) => Ok(Some(true)),
Ok(ProbeVerdict::Pass) => Ok(Some(false)),
Ok(ProbeVerdict::Inconclusive(_)) | Err(_) => Ok::<_, anyhow::Error>(None),
}
};
match fails(Vec::new()).await? {
Some(false) => {}
_ => return Ok(None),
}
let mut set: Vec<usize> = (0..surface.flat.len()).collect();
while set.len() > 1 {
let (a, b) = set.split_at(set.len() / 2);
let (a, b) = (a.to_vec(), b.to_vec());
match fails(a.clone()).await? {
Some(true) => {
set = a;
continue;
}
Some(false) => {}
None => return Ok(None),
}
match fails(b.clone()).await? {
Some(true) => set = b,
_ => return Ok(None),
}
}
Ok(set.first().copied())
}
pub async fn execute(global: &GlobalOpts, args: Args) -> Result<()> {
let store = LearningStore::open(LearningStore::default_root()?)?;
let rules_block = store.rules_prompt_block()?;
let Some(rules_block) = rules_block else {
println!("no rules to validate — run `mecha learn` first");
return Ok(());
};
let wanted_triggers: Vec<&str> = if args.trigger.is_empty() {
vec![
Trigger::Steer.as_str(),
Trigger::Denial.as_str(),
Trigger::Followup.as_str(),
]
} else {
args.trigger.iter().map(String::as_str).collect()
};
let reflexions: Vec<_> = store
.reflexions()?
.into_iter()
.filter(|r| wanted_triggers.contains(&r.trigger.as_str()))
.filter(|r| !args.unprocessed_only || !r.is_processed)
.collect();
if reflexions.is_empty() {
println!("no reflections to probe");
return Ok(());
}
let cwd = std::env::current_dir().context("cannot determine the working directory")?;
let cfg = Config::load(&cwd)?;
let (provider_name, provider_cfg) = cfg.provider(global.provider.as_deref())?;
let provider = mecha_core::provider::build(provider_cfg)?;
let model = global
.model
.clone()
.or_else(|| provider_cfg.model.clone())
.unwrap_or_else(|| provider.default_model().to_string());
let (judge_name, judge_cfg) = cfg.provider(args.judge_provider.as_deref())?;
let judge = Judge::new(
mecha_core::provider::build(judge_cfg)?,
args.judge_model.clone().or_else(|| judge_cfg.model.clone()),
)
.with_max_tokens(16384);
eprintln!(
"probing {} reflection(s) with {model} ({provider_name}), judged by {} ({judge_name})",
reflexions.len(),
judge.model()
);
let needs_replay = reflexions
.iter()
.any(|r| r.trigger != Trigger::Followup.as_str());
let prepared = if needs_replay {
Some(setup::prepare(&global.clone(), false).await?)
} else {
None
};
let sessions_dir = Session::default_dir()?;
let mut improved = 0u32;
let mut regressed = 0u32;
let mut unchanged = 0u32;
let mut inconclusive = 0u32;
let mut skipped = 0u32;
let mut recorded_rows = 0u32;
let surface = RuleSurface::load(&store)?;
let block_hash = rules_hash(&rules_block);
let ledger_rule_ids = surface.rule_ids();
let mut record = |r: &mecha_core::learning::Reflexion,
outcome: &str,
attributed: Option<String>|
-> Result<()> {
store.append_validation(&ValidationRecord {
reflexion_id: r.id.clone(),
trigger: r.trigger.clone(),
domain: r.domain.clone(),
rules_hash: block_hash.clone(),
rule_ids: ledger_rule_ids.clone(),
outcome: outcome.into(),
attributed_rule_id: attributed,
model: model.clone(),
created_at: chrono::Utc::now().to_rfc3339(),
})?;
recorded_rows += 1;
Ok(())
};
for r in &reflexions {
let path = match Session::find(&sessions_dir, &r.session_id) {
Ok(p) => p,
Err(_) => {
eprintln!("· {}: session {} not found; skipping", r.id, r.session_id);
skipped += 1;
continue;
}
};
let (_, convo) = Session::load(&path)?;
let base_system = Session::run_configs(&path)?
.first()
.and_then(|rc| rc.system_prompt.clone())
.map(|s| strip_rules_block(&s))
.unwrap_or_default();
let with_rules = if base_system.is_empty() {
rules_block.clone()
} else {
format!("{base_system}\n\n{rules_block}")
};
if r.trigger != Trigger::Followup.as_str() {
let prepared = prepared.as_ref().expect("built because needs_replay");
let prep = match probe::prepare_probe(&sessions_dir, r)? {
Ok(prep) => prep,
Err(why) => {
eprintln!("· {}: {why}; skipping", r.id);
skipped += 1;
continue;
}
};
let mut arms = Vec::new();
for block in [None, Some(rules_block.as_str())] {
match probe::drive_arm(prepared, provider_cfg, &model, &prep, block).await? {
Ok(v) => arms.push(v),
Err(why) => {
eprintln!("· {}: {why}; skipping", r.id);
break;
}
}
}
let [baseline, with] = &arms[..] else {
skipped += 1;
continue;
};
match probe::compare(
baseline,
with,
&mut improved,
&mut regressed,
&mut unchanged,
&mut inconclusive,
) {
Some(label) => {
println!("· {} [{}, {label}] {}", r.id, r.trigger, r.reflexion_text)
}
None => {
let why = [baseline, with]
.iter()
.find_map(|v| match v {
ProbeVerdict::Inconclusive(w) => Some(w.clone()),
_ => None,
})
.unwrap_or_default();
println!("· {} [{}] inconclusive: {why}", r.id, r.trigger);
}
}
let mut attributed = None;
if matches!((baseline, with), (ProbeVerdict::Pass, ProbeVerdict::Fail))
&& !args.no_attribute
{
match attribute_regression(prepared, provider_cfg, &model, &prep, &surface).await? {
Some(i) => {
let (domain, rule) = &surface.flat[i];
match &rule.id {
Some(id) => {
attributed = Some(id.clone());
println!(" attributed to [{domain}] {}", rule.text);
}
None => println!(
" attributed to a pre-identity rule [{domain}]: {}",
rule.text
),
}
}
None => println!(" no single rule attributable"),
}
}
record(r, outcome_str(baseline, with), attributed)?;
continue;
}
let Some(idx) = locate_followup(&convo.messages, &r.intervention) else {
eprintln!(
"· {}: could not locate the intervention turn; skipping",
r.id
);
skipped += 1;
continue;
};
let mut messages: Vec<Message> = convo.messages[..idx].to_vec();
messages.push(Message::user(r.intervention.clone()));
let mut answers = Vec::new();
for system in [&base_system, &with_rules] {
let request = CompletionRequest {
model: model.clone(),
system: (!system.is_empty()).then(|| system.clone()),
messages: messages.clone(),
tools: Vec::new(),
max_tokens: 4096,
effort: None,
thinking: false,
cache_prompt: true,
};
let response = provider.complete(&request, None).await?;
answers.push(response.message.text());
}
let rubric = format!(
"the answer does what the user's message asks, in the light of this \
known expectation: {}",
r.reflexion_text
);
let mut verdicts = Vec::new();
for answer in &answers {
match judge.assess(&r.intervention, &rubric, answer).await {
Ok(v) => verdicts.push(v.pass),
Err(e) => {
eprintln!("· {}: judge failed ({e:#}); skipping", r.id);
verdicts.clear();
break;
}
}
}
let [baseline, with] = verdicts[..] else {
skipped += 1;
continue;
};
let (label, outcome) = match (baseline, with) {
(false, true) => {
improved += 1;
("IMPROVED", "improved")
}
(true, false) => {
regressed += 1;
("REGRESSED", "regressed")
}
(true, true) => {
unchanged += 1;
("unchanged (both pass)", "unchanged_pass")
}
(false, false) => {
unchanged += 1;
("unchanged (both fail)", "unchanged_fail")
}
};
println!("· {} [{}] {}", r.id, label, r.reflexion_text);
if baseline != with {
println!(" baseline: {}", first_line(&answers[0]));
println!(" with rules: {}", first_line(&answers[1]));
}
record(r, outcome, None)?;
}
println!(
"\n{improved} improved, {regressed} regressed, {unchanged} unchanged, \
{inconclusive} inconclusive, {skipped} skipped (n={}; steers and denials are \
trace-graded, followups judge-graded — read before believing a single flip)",
reflexions.len()
);
if recorded_rows > 0 {
println!(
"{recorded_rows} row(s) appended to the validation ledger — `mecha rules` folds them"
);
store.commit(&format!("validate: {recorded_rows} probe(s) → ledger"));
}
Ok(())
}
fn outcome_str(baseline: &ProbeVerdict, with: &ProbeVerdict) -> &'static str {
match (baseline, with) {
(ProbeVerdict::Inconclusive(_), _) | (_, ProbeVerdict::Inconclusive(_)) => "inconclusive",
(ProbeVerdict::Fail, ProbeVerdict::Pass) => "improved",
(ProbeVerdict::Pass, ProbeVerdict::Fail) => "regressed",
(ProbeVerdict::Pass, _) => "unchanged_pass",
(ProbeVerdict::Fail, _) => "unchanged_fail",
}
}
fn first_line(s: &str) -> String {
let line = s
.lines()
.find(|l| !l.trim().is_empty())
.unwrap_or("")
.trim();
if line.chars().count() > 140 {
format!("{}…", line.chars().take(140).collect::<String>())
} else {
line.to_string()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn temp_store() -> LearningStore {
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos();
let dir = std::env::temp_dir()
.join("mecha-validate-test")
.join(format!("{}-{nanos}", std::process::id()));
LearningStore::open(dir).unwrap()
}
fn rule(text: &str, id: Option<&str>) -> Rule {
Rule {
text: text.into(),
id: id.map(Into::into),
..Default::default()
}
}
#[test]
fn a_full_selection_renders_exactly_the_deployed_block() {
let store = temp_store();
std::fs::write(
store.root().join("rules/behavior.user.toml"),
"[[rules]]\ntext = \"User rule.\"\n",
)
.unwrap();
store
.write_learned_rules(
"behavior",
&[
rule("Learned A.", Some("r-a")),
rule("Learned B.", Some("r-b")),
],
)
.unwrap();
store
.write_learned_rules("writing", &[rule("Sign off briefly.", Some("r-c"))])
.unwrap();
let surface = RuleSurface::load(&store).unwrap();
assert_eq!(surface.flat.len(), 3);
assert_eq!(surface.rule_ids(), vec!["r-a", "r-b", "r-c"]);
let all: Vec<usize> = (0..surface.flat.len()).collect();
assert_eq!(
surface.block_with(&all).unwrap(),
store.rules_prompt_block().unwrap().unwrap()
);
let none = surface.block_with(&[]).unwrap();
assert!(none.contains("User rule."));
assert!(!none.contains("Learned A.") && !none.contains("Sign off"));
let one = surface.block_with(&[1]).unwrap();
assert!(one.contains("Learned B.") && !one.contains("Learned A."));
std::fs::remove_dir_all(store.root()).ok();
}
#[test]
fn retired_and_preidentity_rules_shape_the_surface_correctly() {
let store = temp_store();
let retired = Rule {
text: "Was harmful.".into(),
enabled: false,
id: Some("r-old".into()),
retired_at: Some("2026-08-05T00:00:00Z".into()),
..Default::default()
};
store
.write_learned_rules("behavior", &[rule("No id yet.", None), retired])
.unwrap();
let surface = RuleSurface::load(&store).unwrap();
assert_eq!(
surface.flat.len(),
1,
"retired rules are not on the surface"
);
assert!(surface.rule_ids().is_empty());
assert!(surface.block_with(&[0]).unwrap().contains("No id yet."));
std::fs::remove_dir_all(store.root()).ok();
}
#[test]
fn the_ledger_outcome_vocabulary_covers_the_verdict_grid() {
use ProbeVerdict::*;
let inc = || Inconclusive("why".into());
assert_eq!(outcome_str(&Fail, &Pass), "improved");
assert_eq!(outcome_str(&Pass, &Fail), "regressed");
assert_eq!(outcome_str(&Pass, &Pass), "unchanged_pass");
assert_eq!(outcome_str(&Fail, &Fail), "unchanged_fail");
assert_eq!(outcome_str(&inc(), &Pass), "inconclusive");
assert_eq!(outcome_str(&Pass, &inc()), "inconclusive");
}
}