use crate::analyzer::{AnalyzeCtx, Analyzer, OutcomeInput};
use crate::cal;
use crate::config::{AppliedRecord, LoopPersisted};
use crate::error::{Error, Result};
use crate::manifest::{AnalyzerManifest, Capability};
use crate::model::{normalize_ident, ActionKind, GrainRecord, Origin, Severity, TargetRef};
use crate::recommendation::{
dedup_key, AuditRecord, ObserverType, Proposal, RecStatus, Recommendation, Summary,
MAX_BECAUSE, MAX_EVIDENCE, Checkpoint};
use crate::substrate::{Capabilities, OmsSubstrate, ReadOpts, SubstrateRead};
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use std::collections::{BTreeMap, BTreeSet};
pub const LOOP_NS: &str = "areev-loop";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Scope {
Read,
Write,
Review,
Apply,
Admin,
}
#[derive(Debug, Clone, Default)]
pub struct ScopeSet(Vec<Scope>);
impl ScopeSet {
pub fn of(scopes: &[Scope]) -> Self {
ScopeSet(scopes.to_vec())
}
pub fn all() -> Self {
ScopeSet(vec![Scope::Admin])
}
pub fn has(&self, s: Scope) -> bool {
self.0.contains(&Scope::Admin) || self.0.contains(&s)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Decision {
Approve,
Reject,
}
#[derive(Debug, Clone, Default)]
pub struct RunOptions {
pub min_new: Option<u64>,
pub min_new_errors: Option<u64>,
pub if_stale_ms: Option<i64>,
pub namespaces: Vec<String>,
pub full_sweep: bool,
pub triggering_actor: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RunOutcome {
Ran,
Skipped,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SkipReason {
MinNewNotMet,
NotStale,
LockHeld,
CadenceNotDue,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AnalyzerSkip {
pub id: String,
pub reason: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunResult {
pub outcome: RunOutcome,
#[serde(skip_serializing_if = "Option::is_none")]
pub skip_reason: Option<SkipReason>,
pub new_grains: u64,
pub new_error_events: u64,
pub proposed: u64,
pub deduped: u64,
pub stored: u64,
#[serde(default)]
pub auto_applied: u64,
#[serde(default)]
pub analyzers_run: Vec<String>,
#[serde(default)]
pub analyzers_skipped: Vec<AnalyzerSkip>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub llm_funnel: Option<LlmFunnel>,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct LlmFunnel {
pub evidence: u64,
pub proposed: u64,
pub cited: u64,
pub dropped_uncited: u64,
pub dropped_target: u64,
pub grounded: u64,
pub ground_verdicts: u64,
pub ground_call_failed: bool,
pub kept: u64,
pub stored: u64,
#[serde(default, skip_serializing_if = "is_zero")]
pub advisory_thin_evidence: u64,
}
fn is_zero(n: &u64) -> bool {
*n == 0
}
impl RunResult {
fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
RunResult {
outcome: RunOutcome::Skipped,
skip_reason: Some(reason),
new_grains,
new_error_events,
proposed: 0,
deduped: 0,
stored: 0,
auto_applied: 0,
llm_funnel: None,
analyzers_run: vec![],
analyzers_skipped: vec![],
}
}
pub fn ran(&self) -> bool {
self.outcome == RunOutcome::Ran
}
}
pub struct Engine {
analyzers: Vec<Box<dyn Analyzer>>,
policy: crate::policy::Policy,
llm: Option<Box<dyn crate::llm::LlmBackend>>,
ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
}
struct AnalysisPass {
survivors: Vec<Recommendation>,
proposed: u64,
deduped: u64,
analyzers_run: Vec<String>,
analyzers_skipped: Vec<AnalyzerSkip>,
llm_funnel: Option<LlmFunnel>,
}
impl Engine {
pub fn with_builtins() -> Self {
Engine {
analyzers: crate::analyzer::builtin_analyzers(),
policy: crate::policy::Policy::default(),
llm: None,
ground_llm: None,
}
}
pub fn empty() -> Self {
Engine {
analyzers: vec![],
policy: crate::policy::Policy::default(),
llm: None,
ground_llm: None,
}
}
pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
self.policy = policy;
self
}
pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
self.llm = Some(backend);
self
}
pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
self.ground_llm = Some(backend);
self
}
pub fn policy(&self) -> &crate::policy::Policy {
&self.policy
}
pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
self.analyzers.push(analyzer);
}
pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
&self.analyzers
}
pub fn analyze_only<S: OmsSubstrate>(
&self,
sub: &S,
opts: &RunOptions,
overrides: &BTreeMap<String, Map<String, Value>>,
now_ms: i64,
) -> Result<Vec<Recommendation>> {
let persisted = LoopPersisted::from_value(sub.load_state()?)?;
let analysis_watermark = if opts.full_sweep {
None
} else {
persisted.state.watermark_ms
};
Ok(self
.analysis_pass(
sub,
&persisted,
opts,
overrides,
analysis_watermark,
now_ms,
&[],
)?
.survivors)
}
pub fn run<S: OmsSubstrate>(
&self,
sub: &mut S,
opts: &RunOptions,
now_ms: i64,
) -> Result<RunResult> {
let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
let watermark = persisted.state.watermark_ms;
let analysis_watermark = if opts.full_sweep { None } else { watermark };
let new = count_new(sub, watermark)?;
let (new_grains, new_error_events) = (new.grains, new.error_events);
if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
return Ok(RunResult::skipped(reason, new_grains, new_error_events));
}
let flags_set = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
if !flags_set && !opts.full_sweep {
if let Some(reason) = cadence_gate(&self.policy.cadence, &persisted, new, now_ms) {
return Ok(RunResult::skipped(reason, new_grains, new_error_events));
}
}
let mut outcome_inputs = measure_outcomes(sub, &mut persisted, now_ms)?;
if self.policy.premise_drift {
outcome_inputs.extend(detect_premise_drift(sub, &mut persisted, now_ms)?);
}
let AnalysisPass {
survivors,
proposed,
deduped,
analyzers_run,
analyzers_skipped,
llm_funnel,
} = self.analysis_pass(
&*sub,
&persisted,
opts,
&BTreeMap::new(),
analysis_watermark,
now_ms,
&outcome_inputs,
)?;
let mut stored = 0u64;
let mut auto_applied = 0u64;
for mut rec in survivors {
let spec = rec.to_grain_spec(LOOP_NS)?;
let hash = sub.put_grain(&spec)?;
rec.hash = hash.clone();
let actor = format!("engine:{}", rec.analyzer);
let audit = AuditRecord {
rec_hash: hash.clone(),
from: None,
to: RecStatus::Pending,
actor: actor.clone(),
observer_type: ObserverType::System,
because: "analyzer proposed".into(),
previous_audit_hash: None,
gating: None,
at_ms: now_ms,
};
let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
persisted
.status_index
.insert(hash.clone(), RecStatus::Pending);
persisted.creators.insert(hash.clone(), actor);
if !matches!(rec.origin, Origin::Builtin) {
if let Some(trigger) = &opts.triggering_actor {
persisted.co_creators.insert(hash.clone(), trigger.clone());
}
}
persisted.audit_heads.insert(hash.clone(), audit_hash);
stored += 1;
if self.can_auto_apply(&*sub, &rec) {
self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
auto_applied += 1;
}
}
persisted.state.last_run_ms = Some(now_ms);
persisted.state.watermark_ms = Some(now_ms);
sub.store_state(&persisted.to_value()?)?;
Ok(RunResult {
outcome: RunOutcome::Ran,
skip_reason: None,
new_grains,
new_error_events,
proposed,
deduped,
stored,
auto_applied,
analyzers_run,
analyzers_skipped,
llm_funnel,
})
}
#[allow(clippy::too_many_arguments)]
fn analysis_pass<S: OmsSubstrate>(
&self,
sub: &S,
persisted: &LoopPersisted,
opts: &RunOptions,
external_overrides: &BTreeMap<String, Map<String, Value>>,
analysis_watermark: Option<i64>,
now_ms: i64,
outcome_inputs: &[OutcomeInput],
) -> Result<AnalysisPass> {
let existing = existing_dedup_keys(sub, persisted)?;
let mut analyzers_run = Vec::new();
let mut analyzers_skipped = Vec::new();
let mut candidates: Vec<Recommendation> = Vec::new();
let caps = sub.capabilities();
for analyzer in &self.analyzers {
let m = analyzer.manifest();
let cfg = persisted.config.get(&m.id);
let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
if !enabled {
analyzers_skipped.push(AnalyzerSkip {
id: m.id.clone(),
reason: "disabled".into(),
});
continue;
}
if self.policy.denies(m.family()) {
analyzers_skipped.push(AnalyzerSkip {
id: m.id.clone(),
reason: "denied by host policy".into(),
});
continue;
}
if let Some(missing) = missing_capability(m, caps) {
analyzers_skipped.push(AnalyzerSkip {
id: m.id.clone(),
reason: format!("missing capability: {missing}"),
});
continue;
}
let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
if let Some(extra) = external_overrides.get(&m.id) {
for (key, value) in extra {
param_overrides.insert(key.clone(), value.clone());
}
}
let params = match m.resolve_params(¶m_overrides) {
Ok(p) => p,
Err(e) => {
analyzers_skipped.push(AnalyzerSkip {
id: m.id.clone(),
reason: e.to_string(),
});
continue;
}
};
let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
let ns_slice: &[String] = if ns_owned.is_empty() {
&opts.namespaces
} else {
&ns_owned
};
let reader: &dyn SubstrateRead = sub;
let ctx = AnalyzeCtx::new(
reader,
¶ms,
ns_slice,
analysis_watermark,
now_ms,
outcome_inputs,
);
match analyzer.analyze(&ctx) {
Ok(drafts) => {
analyzers_run.push(m.id.clone());
for draft in drafts {
match stamp(m, ¶ms, draft, now_ms) {
Ok(rec) => candidates.push(rec),
Err(e) => analyzers_skipped.push(AnalyzerSkip {
id: m.id.clone(),
reason: e.to_string(),
}),
}
}
}
Err(e) => analyzers_skipped.push(AnalyzerSkip {
id: m.id.clone(),
reason: e.to_string(),
}),
}
}
let mut funnel = LlmFunnel::default();
if self.llm.is_some() {
candidates.extend(self.discover(
sub,
&candidates,
analysis_watermark,
&opts.namespaces,
now_ms,
&mut funnel,
));
}
let proposed = candidates.len() as u64;
let mut seen = BTreeSet::new();
let mut survivors = Vec::new();
for candidate in candidates {
let family = crate::manifest::analyzer_family(&candidate.analyzer);
let floor = [
severity_floor_for(persisted, &candidate.analyzer),
self.policy.severity_floor(family),
]
.into_iter()
.flatten()
.max();
if floor.is_some_and(|floor| candidate.severity < floor) {
continue;
}
if !seen.insert(candidate.dedup_key.clone()) {
continue;
}
if existing.contains(&candidate.dedup_key) {
continue;
}
if persisted
.cooldowns
.get(&candidate.dedup_key)
.is_some_and(|until| now_ms < *until)
{
continue;
}
survivors.push(candidate);
}
let deduped = proposed - survivors.len() as u64;
if self.llm.is_some() {
self.enrich(&mut survivors);
}
Ok(AnalysisPass {
survivors,
proposed,
deduped,
analyzers_run,
analyzers_skipped,
llm_funnel: self.llm.is_some().then_some(funnel),
})
}
fn discover<S: OmsSubstrate>(
&self,
sub: &S,
candidates: &[Recommendation],
watermark: Option<i64>,
namespaces: &[String],
now_ms: i64,
funnel: &mut LlmFunnel,
) -> Vec<Recommendation> {
let Some(llm) = &self.llm else {
return Vec::new();
};
let findings: Vec<crate::llm::FindingBrief> = candidates
.iter()
.take(32)
.map(|c| crate::llm::FindingBrief {
analyzer: c.analyzer.clone(),
summary: c.summary.render(),
target: c.target_ref.clone(),
severity: c.severity.as_str().to_string(),
})
.collect();
let attribution = self.policy.evidence_attribution;
let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
let mut bundle: BTreeSet<String> = BTreeSet::new();
let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
'cited: for c in candidates {
for h in &c.evidence {
if evidence.len() >= CITED_SEED_CAP {
break 'cited;
}
if !bundle.contains(h) {
if let Ok(Some(g)) = sub.grain(h) {
push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
}
}
}
}
let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
vec![None]
} else {
namespaces.iter().map(|n| Some(n.as_str())).collect()
};
let opts = ReadOpts { live_only: true, since_ms: watermark };
let mut tool_seeded = 0usize;
let seed_successes = self.policy.skills.enabled || self.policy.plans.enabled;
'tools: for want_error in [true, false] {
if !want_error && !seed_successes {
break;
}
for ns in &scan_ns {
if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
for g in recent {
if tool_seeded >= TOOL_SEED_CAP
|| evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
{
break 'tools;
}
if g.is_error() != want_error {
continue;
}
let before = evidence.len();
push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
if evidence.len() > before {
tool_seeded += 1;
}
}
}
}
}
if let Ok(rows) =
sub.grains_of_type(crate::model::grain_type::OBSERVATION, Some(crate::eval::HARNESS_NS), opts)
{
let mut seeded = 0usize;
for g in rows {
if seeded >= HARNESS_SEED_CAP || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE {
break;
}
let kind = g.str_field("observation_kind").unwrap_or_default();
if !HARNESS_EVIDENCE_KINDS.contains(&kind) {
continue;
}
let before = evidence.len();
push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
if evidence.len() > before {
seeded += 1;
}
}
}
'notes: for ns in &scan_ns {
if let Ok(recent) =
sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
{
for g in recent {
if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
break 'notes;
}
push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
}
}
}
'seed: for gt in [
crate::model::grain_type::FACT,
crate::model::grain_type::OBSERVATION,
] {
for ns in &scan_ns {
if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
for g in recent {
if evidence.len() >= EVIDENCE_CAP {
break 'seed;
}
push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
}
}
}
}
funnel.evidence = evidence.len() as u64;
if evidence.is_empty() {
return Vec::new(); }
let (approved, rejected) = self.llm_history(sub);
let base = match self.policy.discover_objective {
crate::policy::DiscoverObjective::ReviewQueue => DISCOVER_INSTRUCTIONS,
crate::policy::DiscoverObjective::Learner => DISCOVER_LEARNER_INSTRUCTIONS,
};
let mut instructions = base.to_string();
if self.policy.skills.enabled {
instructions.push_str(&skill_instructions(self.policy.skills.min_steps));
}
if self.policy.plans.enabled {
instructions.push_str(&plan_instructions(self.policy.plans.min_nodes));
}
let request = crate::llm::LlmRequest {
loop_proto: 1,
op: "discover",
instructions: &instructions,
findings: findings.clone(),
evidence: evidence.clone(),
rejected,
approved,
};
let Ok(body) = serde_json::to_string(&request) else {
return Vec::new();
};
let raw = match llm.complete(&body) {
Ok(r) => r,
Err(_) => return Vec::new(), };
let caps = sub.capabilities();
let mut validated: Vec<ValidatedDraft> = Vec::new();
let drafts: Vec<_> = crate::llm::parse_discover(&raw)
.recommendations
.into_iter()
.take(crate::llm::MAX_LLM_DRAFTS)
.collect();
funnel.proposed = drafts.len() as u64;
let id_to_hash: std::collections::BTreeMap<&str, &str> = evidence
.iter()
.map(|e| (e.id.as_str(), e.hash.as_str()))
.collect();
for d in drafts {
let mut cited: Vec<String> = Vec::new();
for c in &d.evidence {
if let Some(h) = resolve_citation(c, &bundle, &id_to_hash) {
if !cited.contains(&h) {
cited.push(h);
}
}
}
if cited.is_empty() {
funnel.dropped_uncited += 1;
continue; }
let Ok(target) = TargetRef::parse(&d.target) else {
funnel.dropped_target += 1;
continue;
};
let tc = target.target_class();
if !matches!(tc, "memory" | "query" | "code") {
funnel.dropped_target += 1;
continue;
}
let thin = cited.len() < self.policy.min_evidence as usize;
if thin {
funnel.advisory_thin_evidence += 1;
}
let resolved = if thin {
None
} else {
resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps, &self.policy)
};
if tc == "code" && resolved.is_none() {
funnel.dropped_target += 1;
continue;
}
validated.push(ValidatedDraft {
draft: d,
target_ref: target.as_string(),
cited,
resolved,
});
}
funnel.cited = validated.len() as u64;
if validated.is_empty() {
return Vec::new();
}
let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
let outcome_metric = self.outcome_metric_template(sub);
self.verify_drafts(&**llm, ground, validated, &evidence, outcome_metric, now_ms, funnel)
}
fn outcome_metric_template<S: OmsSubstrate>(
&self,
sub: &S,
) -> Option<crate::recommendation::MetricSnapshot> {
let e = self.policy.outcome_evalset.as_ref()?;
let run = crate::eval::newest_eval_run(sub, &e.hash, None).ok().flatten()?;
let baseline = crate::eval::run_value(&run, &e.field)?;
let schedule = e.schedule();
let ms_only: Vec<i64> = schedule.iter().filter_map(Checkpoint::as_ms).collect();
let all_ms = ms_only.len() == schedule.len();
let horizons = if ms_only.is_empty() { vec![86_400_000] } else { ms_only };
Some(crate::recommendation::MetricSnapshot {
metric: format!("evalset:{}:{}", e.hash, e.field),
baseline,
unit: e.field.clone(),
n: run.total(),
window: "per-run".into(),
subject: None,
namespace: None,
relation: None,
query: format!(
"RECALL facts WHERE subject = \"evalset:{}\" AND relation = \"mg:eval_run\"",
e.hash
),
review_after_ms: horizons[0],
horizons_ms: if all_ms { horizons } else { Vec::new() },
checkpoints: if all_ms { Vec::new() } else { schedule },
higher_is_better: e.higher_is_better,
})
}
#[allow(clippy::too_many_arguments)]
fn verify_drafts(
&self,
llm: &dyn crate::llm::LlmBackend,
ground: &dyn crate::llm::LlmBackend,
validated: Vec<ValidatedDraft>,
evidence: &[crate::llm::EvidenceItem],
outcome_metric: Option<crate::recommendation::MetricSnapshot>,
now_ms: i64,
funnel: &mut LlmFunnel,
) -> Vec<Recommendation> {
use crate::llm::*;
let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
cited
.iter()
.filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
.collect()
};
let claims: Vec<GroundItem> = validated
.iter()
.enumerate()
.map(|(i, v)| GroundItem {
id: i,
claim: claim_text(&v.draft, v.resolved.as_ref()),
evidence: ev_for(&v.cited),
})
.collect();
let ground_req = GroundRequest {
loop_proto: 1,
op: "ground",
instructions: GROUND_INSTRUCTIONS,
claims,
};
let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
.ok()
.and_then(|b| ground.complete(&b).ok())
{
Some(raw) => {
let parsed = parse_ground(&raw);
funnel.ground_verdicts = parsed.results.len() as u64;
parsed
.results
.into_iter()
.filter(|r| r.supported)
.map(|r| r.id)
.collect()
}
None => {
funnel.ground_call_failed = true;
return Vec::new();
}
};
funnel.grounded = grounded.len() as u64;
if grounded.is_empty() {
return Vec::new();
}
let items: Vec<VerifyItem> = validated
.iter()
.enumerate()
.filter(|(i, _)| grounded.contains(i))
.map(|(i, v)| VerifyItem {
id: i,
summary: claim_text(&v.draft, v.resolved.as_ref()),
target: v.target_ref.clone(),
evidence: ev_for(&v.cited),
})
.collect();
let verify_req = VerifyRequest {
loop_proto: 1,
op: "verify",
instructions: VERIFY_INSTRUCTIONS,
findings: items,
};
let verdicts: std::collections::BTreeMap<usize, f64> =
match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
Some(raw) => parse_verify(&raw)
.results
.into_iter()
.filter(|r| r.keep)
.map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
.collect(),
None => return Vec::new(),
};
funnel.kept = verdicts.len() as u64;
let mut out = Vec::new();
for (i, v) in validated.into_iter().enumerate() {
if let Some(&conf) = verdicts.get(&i) {
if conf >= MIN_LLM_CONFIDENCE {
let mut rec = stamp_llm(
llm.model(),
&v.draft,
v.target_ref,
v.cited,
v.resolved,
conf,
now_ms,
);
if rec.rollbackable {
rec.metric = outcome_metric.clone();
}
out.push(rec);
}
}
}
funnel.stored = out.len() as u64;
out
}
fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
const MAX: usize = 20;
let Ok(mut recs) = self.recommendations(sub, None) else {
return (Vec::new(), Vec::new());
};
recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
let mut approved = Vec::new();
let mut rejected = Vec::new();
for r in &recs {
match r.status {
RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
if approved.len() < MAX =>
{
approved.push(r.summary.render());
}
RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
_ => {}
}
}
(approved, rejected)
}
fn enrich(&self, survivors: &mut [Recommendation]) {
let Some(llm) = &self.llm else {
return;
};
if survivors.is_empty() {
return;
}
let findings: Vec<crate::llm::FindingBrief> = survivors
.iter()
.map(|r| crate::llm::FindingBrief {
analyzer: r.analyzer.clone(),
summary: r.summary.render(),
target: r.target_ref.clone(),
severity: r.severity.as_str().to_string(),
})
.collect();
let request = crate::llm::LlmRequest {
loop_proto: 1,
op: "enrich",
instructions: ENRICH_INSTRUCTIONS,
findings,
evidence: Vec::new(),
rejected: Vec::new(),
approved: Vec::new(),
};
let Ok(body) = serde_json::to_string(&request) else {
return;
};
let raw = match llm.complete(&body) {
Ok(r) => r,
Err(_) => return,
};
for note in crate::llm::parse_enrich(&raw).notes {
if note.guidance.trim().is_empty() {
continue;
}
if let Some(r) = survivors
.iter_mut()
.find(|r| r.target_ref == note.target && r.guidance.is_none())
{
r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
}
}
}
fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
if !rec.origin.auto_apply_eligible() || rec.destructive {
return false;
}
let manifest_ok = self
.analyzers
.iter()
.map(|a| a.manifest())
.find(|m| m.id == rec.analyzer)
.is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
if !manifest_ok {
return false;
}
let Ok(target) = TargetRef::parse(&rec.target_ref) else {
return false;
};
let family = crate::manifest::analyzer_family(&rec.analyzer);
if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
return false;
}
match &rec.proposal {
Proposal::Cal { cal } => cal
.lines()
.map(str::trim)
.filter(|l| !l.is_empty())
.all(|l| supersede_is_value_identical(sub, l)),
_ => false,
}
}
fn auto_apply<S: OmsSubstrate>(
&self,
sub: &mut S,
p: &mut LoopPersisted,
rec: &Recommendation,
now_ms: i64,
) -> Result<()> {
let mut created = Vec::new();
if let Proposal::Cal { cal } = &rec.proposal {
if cal.lines().map(str::trim).any(is_definition_statement) {
return Err(Error::InvalidProposal(
"a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
auto-applied: it changes what every future context contains, so it \
requires a human APPROVE + APPLY with BECAUSE"
.into(),
));
}
for r in sub.execute_cal(cal)? {
if let Some(h) = r.get("hash").and_then(Value::as_str) {
created.push(h.to_string());
}
}
}
let applied = AppliedRecord {
applied_at_ms: now_ms,
target_ref: rec.target_ref.clone(),
rollbackable: rec.rollbackable,
created_hashes: created,
inverse_cal: None,
metric: rec.metric.clone(),
};
let prev = p.audit_heads.get(&rec.hash).cloned();
let audit = AuditRecord {
rec_hash: rec.hash.clone(),
from: Some(RecStatus::Pending),
to: RecStatus::Applied,
actor: "policy:auto".into(),
observer_type: ObserverType::Policy,
because: "auto-applied per host policy".into(),
previous_audit_hash: prev,
gating: None,
at_ms: now_ms,
};
let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
p.audit_heads.insert(rec.hash.clone(), audit_hash);
p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
p.applied.insert(rec.hash.clone(), applied);
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub fn review<S: OmsSubstrate>(
&self,
sub: &mut S,
rec_hash: &str,
decision: Decision,
actor: &str,
observer: ObserverType,
scopes: &ScopeSet,
because: &str,
now_ms: i64,
) -> Result<()> {
if !scopes.has(Scope::Review) {
return Err(Error::ScopeDenied("review".into()));
}
let because = validate_because(because)?;
let mut p = LoopPersisted::from_value(sub.load_state()?)?;
let status = *p
.status_index
.get(rec_hash)
.ok_or_else(|| Error::NotFound(rec_hash.into()))?;
let to = match decision {
Decision::Approve => RecStatus::Approved,
Decision::Reject => RecStatus::Rejected,
};
if !status.can_transition_to(to, false) {
return Err(Error::LifecycleViolation(format!(
"{} -> {}",
status.as_str(),
to.as_str()
)));
}
if to == RecStatus::Approved {
if let Some(creator) = p.creators.get(rec_hash) {
if creator == actor {
return Err(Error::SelfApproval(format!(
"{actor} created this recommendation"
)));
}
}
if let Some(trigger) = p.co_creators.get(rec_hash) {
if trigger == actor {
return Err(Error::SelfApproval(format!(
"{actor} triggered the run that authored this recommendation"
)));
}
}
}
let prev = p.audit_heads.get(rec_hash).cloned();
let audit = AuditRecord {
rec_hash: rec_hash.into(),
from: Some(status),
to,
actor: actor.into(),
observer_type: observer,
because,
previous_audit_hash: prev,
gating: None,
at_ms: now_ms,
};
let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
p.audit_heads.insert(rec_hash.into(), audit_hash);
p.status_index.insert(rec_hash.into(), to);
if to == RecStatus::Rejected {
if let Ok(rec) = load_rec(sub, rec_hash) {
strike_cooldown(&mut p, rec.dedup_key, now_ms);
}
}
sub.store_state(&p.to_value()?)?;
Ok(())
}
pub fn preflight_apply<S: OmsSubstrate>(
&self,
sub: &S,
rec_hash: &str,
scopes: &ScopeSet,
allow_destructive: bool,
has_gating: bool,
) -> Result<()> {
if !scopes.has(Scope::Apply) {
return Err(Error::ScopeDenied("apply".into()));
}
let rec = load_rec(sub, rec_hash)?;
if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
return Err(Error::DestructiveGated(
"destructive apply requires admin scope + allow_destructive".into(),
));
}
ensure_executable(rec.action_kind, &rec.proposal)?;
if requires_gating(rec.action_kind) && !has_gating {
return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub fn apply<S: OmsSubstrate>(
&self,
sub: &mut S,
rec_hash: &str,
actor: &str,
observer: ObserverType,
scopes: &ScopeSet,
because: &str,
allow_destructive: bool,
now_ms: i64,
) -> Result<AppliedRecord> {
self.apply_inner(
sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
)
}
pub fn gating_evidence<S: OmsSubstrate>(
&self,
sub: &S,
rec_hash: &str,
run_id: &str,
) -> Result<crate::recommendation::GatingEvidence> {
let rec = self
.recommendations(sub, None)?
.into_iter()
.find(|r| r.hash == rec_hash)
.ok_or_else(|| {
Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
})?;
let pin = rec.evalset_hash.ok_or_else(|| {
Error::InvalidProposal(
"this recommendation pins no evalset — a gating run applies only \
to code and adapter revisions"
.into(),
)
})?;
match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
Some(run) => Ok(crate::recommendation::GatingEvidence {
evalset_hash: pin,
run_id: run.run_id,
passed: run.passed,
failed: run.failed,
}),
None => Err(Error::InvalidProposal(format!(
"no recorded gate run '{run_id}' for evalset {pin} — run \
`areev eval run --evalset {pin} ...` first"
))),
}
}
#[allow(clippy::too_many_arguments)]
pub fn apply_gated<S: OmsSubstrate>(
&self,
sub: &mut S,
rec_hash: &str,
actor: &str,
observer: ObserverType,
scopes: &ScopeSet,
because: &str,
allow_destructive: bool,
gating: &crate::recommendation::GatingEvidence,
now_ms: i64,
) -> Result<AppliedRecord> {
self.apply_inner(
sub,
rec_hash,
actor,
observer,
scopes,
because,
allow_destructive,
Some(gating),
now_ms,
)
}
#[allow(clippy::too_many_arguments)]
fn apply_inner<S: OmsSubstrate>(
&self,
sub: &mut S,
rec_hash: &str,
actor: &str,
observer: ObserverType,
scopes: &ScopeSet,
because: &str,
allow_destructive: bool,
gating: Option<&crate::recommendation::GatingEvidence>,
now_ms: i64,
) -> Result<AppliedRecord> {
if !scopes.has(Scope::Apply) {
return Err(Error::ScopeDenied("apply".into()));
}
let because = validate_because(because)?;
let mut p = LoopPersisted::from_value(sub.load_state()?)?;
let status = *p
.status_index
.get(rec_hash)
.ok_or_else(|| Error::NotFound(rec_hash.into()))?;
if !status.can_transition_to(RecStatus::Applied, false) {
return Err(Error::LifecycleViolation(format!(
"{} -> applied (approve first)",
status.as_str()
)));
}
let rec = load_rec(sub, rec_hash)?;
if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
return Err(Error::DestructiveGated(
"destructive apply requires admin scope + allow_destructive".into(),
));
}
if requires_gating(rec.action_kind) {
let g = gating
.ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
let pin = rec.evalset_hash.as_deref().unwrap_or_default();
if g.evalset_hash != pin {
return Err(Error::InvalidProposal(format!(
"gating ran evalset {} but the recommendation is pinned \
to {pin} (Rule E1)",
g.evalset_hash
)));
}
match sub.grain(pin)? {
Some(evalset) if evalset.is_live() => {}
Some(_) => {
return Err(Error::InvalidProposal(
"the pinned evalset was superseded after gating — \
the recommendation must re-gate (Rule E1)"
.into(),
))
}
None => {
return Err(Error::InvalidProposal(format!(
"pinned evalset {pin} not found in the substrate"
)))
}
}
if g.failed > 0 {
return Err(Error::InvalidProposal(format!(
"the gating run failed {}/{} cases — a failing gate \
admits nothing",
g.failed,
g.passed + g.failed
)));
}
}
let mut created = Vec::new();
let mut inverse_cal: Option<String> = None;
match &rec.proposal {
Proposal::Cal { cal } => {
for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
if !is_definition_statement(line) {
continue;
}
match sub.definition_inverse(line)? {
Some(inv) => inverse_cal = Some(inv),
None => {
return Err(Error::InvalidProposal(format!(
"this substrate cannot record a rollback inverse for {line:?}; \
a definition rewrite that ROLLBACK could not undo is refused \
rather than applied"
)))
}
}
}
let rows = sub.execute_cal(cal)?;
for r in rows {
if let Some(h) = r.get("hash").and_then(Value::as_str) {
created.push(h.to_string());
}
}
}
Proposal::Edit { .. } => {
return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
}
Proposal::Data { data } if requires_gating(rec.action_kind) => {
let relation = if rec.action_kind == ActionKind::AdapterRevision {
"mg:adapter_promotion"
} else {
"mg:code_promotion"
};
let mut promoted = data.clone();
if let Some(Value::String(src)) = promoted.remove("source") {
let address = sub.put_blob(src.as_bytes())?;
promoted.insert("code_address".into(), Value::from(address));
}
let mut spec = crate::substrate::GrainSpec::new(
crate::model::grain_type::FACT,
LOOP_NS,
)
.with_field("subject", rec.target_ref.clone())
.with_field("relation", relation)
.with_field(
"object",
serde_json::to_string(&promoted)
.map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
)
.with_field("rec_hash", rec_hash.to_string());
if let Some(g) = gating {
spec = spec
.with_field("gating_evalset", g.evalset_hash.clone())
.with_field("gating_run_id", g.run_id.clone());
}
created.push(sub.put_grain(&spec)?);
}
Proposal::Data { data } => {
let revert_of = data
.get("revert_of")
.and_then(Value::as_str)
.ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
self.rollback(
sub,
revert_of,
actor,
observer,
scopes,
&because,
now_ms,
)?;
p = LoopPersisted::from_value(sub.load_state()?)?;
if let Ok(reverted) = load_rec(sub, revert_of) {
strike_cooldown(&mut p, reverted.dedup_key, now_ms);
}
}
}
let applied = AppliedRecord {
applied_at_ms: now_ms,
target_ref: rec.target_ref.clone(),
rollbackable: rec.rollbackable,
created_hashes: created,
inverse_cal,
metric: rec.metric.clone(),
};
let prev = p.audit_heads.get(rec_hash).cloned();
let audit = AuditRecord {
rec_hash: rec_hash.into(),
from: Some(status),
to: RecStatus::Applied,
actor: actor.into(),
observer_type: observer,
because,
previous_audit_hash: prev,
gating: gating.cloned(),
at_ms: now_ms,
};
let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
p.audit_heads.insert(rec_hash.into(), audit_hash);
p.status_index.insert(rec_hash.into(), RecStatus::Applied);
p.applied.insert(rec_hash.into(), applied.clone());
sub.store_state(&p.to_value()?)?;
Ok(applied)
}
#[allow(clippy::too_many_arguments)]
pub fn rollback<S: OmsSubstrate>(
&self,
sub: &mut S,
rec_hash: &str,
actor: &str,
observer: ObserverType,
scopes: &ScopeSet,
because: &str,
now_ms: i64,
) -> Result<()> {
if !scopes.has(Scope::Apply) {
return Err(Error::ScopeDenied("apply".into()));
}
let because = validate_because(because)?;
let mut p = LoopPersisted::from_value(sub.load_state()?)?;
let status = *p
.status_index
.get(rec_hash)
.ok_or_else(|| Error::NotFound(rec_hash.into()))?;
if !status.can_transition_to(RecStatus::RolledBack, false) {
return Err(Error::LifecycleViolation(format!(
"{} -> rolled_back",
status.as_str()
)));
}
let applied = p
.applied
.get(rec_hash)
.cloned()
.ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
if !applied.rollbackable {
return Err(Error::LifecycleViolation(
"recommendation is non-rollbackable (FORGET has no inverse)".into(),
));
}
for h in &applied.created_hashes {
sub.retract(h, &format!("rollback of {rec_hash}"))?;
}
if let Some(inverse) = &applied.inverse_cal {
sub.execute_cal(inverse)?;
}
let prev = p.audit_heads.get(rec_hash).cloned();
let audit = AuditRecord {
rec_hash: rec_hash.into(),
from: Some(status),
to: RecStatus::RolledBack,
actor: actor.into(),
observer_type: observer,
because,
previous_audit_hash: prev,
gating: None,
at_ms: now_ms,
};
let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
p.audit_heads.insert(rec_hash.into(), audit_hash);
p.status_index
.insert(rec_hash.into(), RecStatus::RolledBack);
sub.store_state(&p.to_value()?)?;
Ok(())
}
pub fn recommendations<S: OmsSubstrate>(
&self,
sub: &S,
status_filter: Option<RecStatus>,
) -> Result<Vec<Recommendation>> {
let p = LoopPersisted::from_value(sub.load_state()?)?;
let grains = sub.grains_of_type(
crate::model::grain_type::RECOMMENDATION,
Some(LOOP_NS),
ReadOpts {
live_only: false,
since_ms: None,
},
)?;
let mut out = Vec::new();
for g in grains {
let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
rec.status = p
.status_index
.get(&g.hash)
.copied()
.unwrap_or(RecStatus::Pending);
if let Some(f) = status_filter {
if rec.status != f {
continue;
}
}
out.push(rec);
}
out.sort_by(|a, b| {
b.severity
.cmp(&a.severity)
.then(a.created_at_ms.cmp(&b.created_at_ms))
.then(a.dedup_key.cmp(&b.dedup_key))
.then(a.hash.cmp(&b.hash))
});
Ok(out)
}
pub fn analyzer_settings<S: OmsSubstrate>(
&self,
sub: &S,
) -> Result<Vec<crate::config::AnalyzerSetting>> {
let p = LoopPersisted::from_value(sub.load_state()?)?;
Ok(self
.analyzers
.iter()
.map(|a| {
let m = a.manifest();
let cfg = p.config.get(&m.id);
crate::config::AnalyzerSetting {
id: m.id.clone(),
title: m.title.clone(),
description: m.description.clone(),
tier: format!("{:?}", m.tier),
trust_class: format!("{:?}", m.trust_class).to_lowercase(),
default_on: m.default_on,
enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
severity_floor: cfg
.and_then(|c| c.severity_floor)
.map(|s| s.as_str().to_string()),
}
})
.collect())
}
pub fn set_analyzer_config<S: OmsSubstrate>(
&self,
sub: &mut S,
analyzer_id: &str,
update: crate::config::AnalyzerConfigUpdate,
scopes: &ScopeSet,
) -> Result<crate::config::AnalyzerConfig> {
if !scopes.has(Scope::Admin) {
return Err(Error::ScopeDenied("admin".into()));
}
let manifest = self
.analyzers
.iter()
.map(|a| a.manifest())
.find(|m| m.id == analyzer_id)
.ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
if let Some(params) = &update.params {
manifest.resolve_params(params)?;
}
let mut p = LoopPersisted::from_value(sub.load_state()?)?;
let cfg = p.config.entry(analyzer_id.to_string()).or_default();
if let Some(enabled) = update.enabled {
cfg.enabled = Some(enabled);
}
if update.clear_floor {
cfg.severity_floor = None;
} else if let Some(floor) = update.severity_floor {
cfg.severity_floor = Some(floor);
}
if let Some(params) = update.params {
cfg.params = params;
}
if let Some(ns) = update.namespaces {
cfg.namespaces = ns;
}
let stored = cfg.clone();
sub.store_state(&p.to_value()?)?;
Ok(stored)
}
pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
let p = LoopPersisted::from_value(sub.load_state()?)?;
let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
out.sort_by(|a, b| {
a.measured_at_ms
.cmp(&b.measured_at_ms)
.then(a.horizon_ms.cmp(&b.horizon_ms))
.then(a.metric.cmp(&b.metric))
.then(a.rec_hash.cmp(&b.rec_hash))
});
Ok(out)
}
pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
let p = LoopPersisted::from_value(sub.load_state()?)?;
let new = count_new(sub, p.state.watermark_ms)?;
let (grains_since_run, error_events_since_run) = (new.grains, new.error_events);
let recs = self.recommendations(sub, None)?;
let mut pending = 0;
let mut applied = 0;
for r in &recs {
match r.status {
RecStatus::Pending => pending += 1,
RecStatus::Applied => applied += 1,
_ => {}
}
}
let stale = match p.state.last_run_ms {
None => true,
Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
};
Ok(Health {
last_run_ms: p.state.last_run_ms,
grains_since_run,
error_events_since_run,
pending,
applied,
total: recs.len() as u64,
stale,
})
}
pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
let recs = self.recommendations(sub, None)?;
let mut m = LlmMetrics::default();
for r in &recs {
if !matches!(r.origin, Origin::Llm { .. }) {
continue;
}
m.proposed += 1;
match r.status {
RecStatus::Pending => m.pending += 1,
RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
RecStatus::Rejected => m.rejected += 1,
RecStatus::Expired => {}
}
}
let decided = m.approved + m.rejected;
m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
Ok(m)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Health {
#[serde(skip_serializing_if = "Option::is_none")]
pub last_run_ms: Option<i64>,
pub grains_since_run: u64,
pub error_events_since_run: u64,
pub pending: u64,
pub applied: u64,
pub total: u64,
pub stale: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct LlmMetrics {
pub proposed: u64,
pub pending: u64,
pub approved: u64,
pub rejected: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub approval_rate: Option<f64>,
}
fn measure_outcomes<S: OmsSubstrate>(
sub: &S,
p: &mut LoopPersisted,
now_ms: i64,
) -> Result<Vec<OutcomeInput>> {
let mut due: Vec<(String, crate::config::AppliedRecord, Checkpoint)> = Vec::new();
for (h, a) in &p.applied {
if p.status_index.get(h) != Some(&RecStatus::Applied) {
continue;
}
let Some(metric) = &a.metric else { continue };
let done = p.measured.get(h).cloned().unwrap_or_default();
for cp in metric.schedule() {
if !done.contains(&cp) && checkpoint_due(sub, metric, a.applied_at_ms, cp, now_ms)? {
due.push((h.clone(), a.clone(), cp));
}
}
}
let mut out = Vec::new();
for (rec_hash, applied, checkpoint) in due {
let metric = applied.metric.as_ref().unwrap();
let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
continue; };
let baseline = baseline_at_apply(sub, metric, applied.applied_at_ms)?;
let regressed = crate::recommendation::is_regression(
baseline,
current,
metric.higher_is_better,
);
p.outcomes.entry(rec_hash.clone()).or_default().push(
crate::recommendation::OutcomeResult {
rec_hash: rec_hash.clone(),
metric: metric.metric.clone(),
baseline,
current,
verdict: if regressed { "regressed" } else { "held" }.into(),
horizon_ms: checkpoint.as_ms().unwrap_or(0),
checkpoint: checkpoint.as_ms().is_none().then_some(checkpoint),
measured_at_ms: now_ms,
},
);
p.measured.entry(rec_hash.clone()).or_default().push(checkpoint);
if regressed {
out.push(OutcomeInput {
rec_hash,
target_ref: applied.target_ref.clone(),
metric: metric.metric.clone(),
baseline,
current,
unit: metric.unit.clone(),
higher_is_better: metric.higher_is_better,
});
}
}
Ok(out)
}
fn checkpoint_due<S: SubstrateRead>(
sub: &S,
metric: &crate::recommendation::MetricSnapshot,
applied_at_ms: i64,
cp: Checkpoint,
now_ms: i64,
) -> Result<bool> {
Ok(match cp {
Checkpoint::AfterMs(ms) => now_ms - applied_at_ms >= ms,
Checkpoint::AfterRuns(n) => match crate::eval::parse_evalset_metric(&metric.metric) {
Some((evalset, _)) => {
crate::eval::eval_runs(sub, evalset, Some(applied_at_ms + 1))?.len() >= n as usize
}
None => false,
},
Checkpoint::AfterGrains(n) => count_new(sub, Some(applied_at_ms))?.grains >= n as u64,
})
}
fn baseline_at_apply<S: SubstrateRead>(
sub: &S,
metric: &crate::recommendation::MetricSnapshot,
applied_at_ms: i64,
) -> Result<f64> {
if let Some((evalset, field)) = crate::eval::parse_evalset_metric(&metric.metric) {
if let Some(run) = crate::eval::newest_eval_run_before(sub, evalset, applied_at_ms)? {
if let Some(v) = crate::eval::run_value(&run, field) {
return Ok(v);
}
}
}
Ok(metric.baseline)
}
pub(crate) fn measure_metric<S: SubstrateRead>(
sub: &S,
metric: &crate::recommendation::MetricSnapshot,
since_ms: i64,
) -> Result<Option<f64>> {
match metric.metric.as_str() {
"tool_error_recurrence" => {
let Some(tool) = &metric.subject else { return Ok(None) };
let tools = sub.grains_of_type(
crate::model::grain_type::TOOL,
None,
ReadOpts { live_only: true, since_ms: Some(since_ms) },
)?;
let n = tools
.iter()
.filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
.filter(|t| {
metric.relation.as_deref().is_none_or(|sig| {
crate::analyzers::tool_failure::normalize_signature(
t.tool_content().unwrap_or(""),
) == sig
})
})
.count();
Ok(Some(n as f64))
}
"contradiction_recurrence" => {
let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
return Ok(None);
};
let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
let distinct: BTreeSet<String> = facts
.iter()
.filter(|f| {
f.fact_relation()
.is_some_and(|r| normalize_ident(r) == *relation)
})
.filter_map(|f| f.fact_object().map(normalize_ident))
.collect();
Ok(Some(distinct.len().saturating_sub(1) as f64))
}
m if m.starts_with("evalset:") => {
let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
return Ok(None);
};
let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
return Ok(None);
};
Ok(crate::eval::run_value(&run, field))
}
_ => Ok(None),
}
}
fn scoped_live_facts<S: SubstrateRead>(
sub: &S,
namespace: Option<&str>,
subject: &str,
) -> Result<Vec<GrainRecord>> {
let facts = sub.grains_of_type(
crate::model::grain_type::FACT,
None,
ReadOpts { live_only: true, since_ms: None },
)?;
Ok(facts
.into_iter()
.filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
.filter(|f| {
f.fact_subject()
.is_some_and(|s| normalize_ident(s) == subject)
})
.collect())
}
fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
return false;
};
if fields.is_empty() {
return false;
}
let Ok(Some(grain)) = sub.grain(&target) else {
return false;
};
if grain.valid_to_ms.is_some() {
return false;
}
fields.iter().all(|(k, v)| {
if k == "namespace" {
return v
.as_str()
.is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
}
match (v, grain.fields.get(k)) {
(Value::String(a), Some(Value::String(b))) => {
normalize_ident(a) == normalize_ident(b)
}
(a, Some(b)) => a == b,
(_, None) => false,
}
})
}
const EVIDENCE_CAP: usize = 64;
const CITED_SEED_CAP: usize = 24;
const TOOL_SEED_CAP: usize = 16;
const NOTE_SEED_CAP: usize = 8;
const HARNESS_EVIDENCE_KINDS: &[&str] = &["fold_summary"];
const HARNESS_SEED_CAP: usize = 6;
const LENS_RESERVE: usize = 24;
const MIN_LLM_CONFIDENCE: f64 = 0.75;
macro_rules! discover_instructions {
($scoring:literal) => {
concat!(
"You review an agent's memory for quality. \
Given deterministic findings and the evidence they cite, propose ADDITIONAL \
findings the deterministic checks would miss (e.g. a semantic contradiction, a \
stale assumption, a duplicated meaning, a recurring preventable mistake, a \
recurring cost or hand-off the agent's own setup could remove). \
The deterministic findings already cover what the ERROR TEXT says; restating \
one of them earns nothing. The evidence may also contain OUTCOME records — a \
run's observable shape together with whether it was accepted or rejected. A \
problem that raised no error at all is exactly the kind the deterministic \
checks cannot see, so compare the rejected outcomes against the accepted \
ones: a feature they share and the accepted ones lack is a candidate rule. \
Require at least two rejected outcomes before proposing one — a single \
rejection is an anecdote, not a pattern. ",
$scoring,
" The 'approved' and 'rejected' lists, when \
present, show findings this reviewer recently accepted or rejected — prefer the \
kind they accept and avoid the kind they reject. Every proposal MUST cite one \
or more evidence items from the bundle by their 'id' (or 'hash'), name a \
'target', and include your confidence 0.0-1.0. Return JSON: \
{\"recommendations\":[{\"summary\":\"...\",\
\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<id>\"],\
\"confidence\":0.0,\"proposal\":{...}}]}. \
OMIT 'proposal' for an advisory finding — one worth a human's attention that \
you are not asking to change anything. Include it ONLY when the evidence \
supports a specific change, choosing exactly one kind: \
(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
an action the agent itself takes on the next occasion. Either ADD an action \
it is failing to take ('Record the vendor name and the amount on every \
invoice, not just the date') or ORDER one that goes wrong ('Refund a \
subscription before cancelling it; refunds on cancelled subscriptions are \
refused'). Name the action, not a check on it: 'validate', 'verify' and \
'ensure ... is correct' describe a review step the agent has no way to \
perform, and such a rule changes nothing even once applied. \
(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
entity target — a durable fact the agent keeps having to be told (an alias, a \
settled default, a preference). 'relation' is a short identifier (letters, \
digits, _ - . :), not a sentence. \
(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
assembles the agent's context, when the evidence shows it retrieves the wrong \
things. Give the FULL new body; it replaces the old one. \
(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
workflow. Only these paths are editable: 'edges.<i>.cond', \
'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
— full replacement source for that tool. It is applied only after a recorded \
evaluation run passes, so propose one only when the evidence shows the current \
code is the defect. \
The subject of a fact, the name of a query, the plan hash and the tool name \
all come from 'target' — do not repeat them inside 'proposal'. A proposal \
becomes a change a human reviewer may apply, so it must be fully supported by \
the cited evidence. Propose nothing you cannot ground in the evidence."
)
};
}
const DISCOVER_INSTRUCTIONS: &str = discover_instructions!(
"SCORING: propose a finding ONLY if you \
are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
earns 0. When in doubt, propose nothing — an empty list is the correct answer \
when there is nothing worth flagging."
);
const DISCOVER_LEARNER_INSTRUCTIONS: &str = discover_instructions!(
"SCORING: you are the learning stage of a deployed agent, and what you \
propose now is what it will do differently next time — a lesson you withhold \
is a mistake it repeats. A correct, actionable proposal earns 1; a wrong or \
trivial one is penalized 1; returning nothing while the evidence holds a \
recurring failure, two or more rejected outcomes, an instruction from a \
person, or a multi-step procedure the agent completed successfully that no \
saved skill or plan covers, is ALSO penalized 1. Abstain only when the \
evidence shows none of those. Prefer the one proposal that addresses the most \
frequent or most costly failure — or, when nothing failed, the procedure that \
worked — over several speculative ones, and report your confidence honestly — \
an independent verifier, not you, decides what survives."
);
const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
the facts it relies on are actually present in the cited evidence, NOT that its \
conclusion is stated verbatim. Decompose the finding into the factual claims it \
depends on. Mark supported=true when those facts are present in the evidence \
(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
'possible' findings with no concrete defect, and reject any claimed \
inconsistency or contradiction that is not backed by at least two actually \
conflicting facts in the cited evidence. (2) Context — does the finding \
correctly read its cited evidence, or misinterpret what the grains say? Keep a \
finding when it names a genuine, specific problem grounded in its evidence and \
materially useful to a human reviewer; otherwise reject it, and default to \
keep=false when uncertain. Do NOT reject a finding for being 'already known', \
redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
grounded cross-fact inconsistency with two conflicting facts is exactly what to \
KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
guidance note to help a human reviewer decide. Do not restate the finding. Return \
JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
fn push_evidence(
evidence: &mut Vec<crate::llm::EvidenceItem>,
bundle: &mut BTreeSet<String>,
ns_by_hash: &mut std::collections::BTreeMap<String, String>,
g: &GrainRecord,
attribution: crate::policy::EvidenceAttribution,
) {
if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
evidence.push(crate::llm::EvidenceItem {
id: format!("e{}", evidence.len() + 1),
hash: g.hash.clone(),
grain_type: g.grain_type.clone(),
text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
});
}
}
pub(crate) fn resolve_citation(
cite: &str,
bundle: &BTreeSet<String>,
id_to_hash: &std::collections::BTreeMap<&str, &str>,
) -> Option<String> {
let cite = cite.trim();
if bundle.contains(cite) {
return Some(cite.to_string());
}
if let Some(h) = id_to_hash.get(cite) {
return Some((*h).to_string());
}
const MIN_PREFIX: usize = 12;
if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
let lower = cite.to_ascii_lowercase();
let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
if let (Some(h), None) = (it.next(), it.next()) {
return Some(h.clone());
}
}
None
}
fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
return format!("{s} {r} {o}");
}
if let Some(t) = g.tool_name() {
let status = if g.is_error() { "error" } else { "ok" };
let out = g.tool_content().unwrap_or("");
let input = match g.fields.get("input") {
Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
_ => String::new(),
};
return format!("tool {t}{input} {status}: {out}");
}
for key in ["content", "body", "text", "summary", "object"] {
if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
if v.is_empty() {
continue;
}
if attribution == crate::policy::EvidenceAttribution::Anonymous {
return v.to_string();
}
if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
if !who.is_empty() {
let kind = g
.fields
.get("observer_type")
.and_then(|v| v.as_str())
.unwrap_or("");
let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
let mut prefix = if kind == "human" {
format!("{who} (a person) said")
} else {
format!("{who} observed")
};
if !about.is_empty() {
prefix.push_str(&format!(" of {about}"));
}
return format!("{prefix}: {v}");
}
}
return v.to_string();
}
}
String::new()
}
fn sanitize_lesson(s: &str) -> String {
sanitize_line(s, crate::llm::MAX_LESSON_LEN)
}
fn sanitize_line(s: &str, max: usize) -> String {
let cleaned: String =
s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
crate::llm::cap(cleaned.trim(), max)
}
fn sanitize_relation(s: &str) -> Option<String> {
let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
if r.is_empty()
|| !r
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
{
return None;
}
Some(r)
}
fn safe_definition_body(body: &str) -> bool {
if body.contains('{') || body.contains('}') {
return false;
}
!body
.split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
.any(|tok| {
["FORGET", "PURGE", "DROP", "DEFINE"]
.iter()
.any(|kw| tok.eq_ignore_ascii_case(kw))
})
}
fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
match resolved {
Some(r) => format!("{summary} {}", r.rendered),
None => summary,
}
}
struct ValidatedDraft {
draft: crate::llm::LlmDraft,
target_ref: String,
cited: Vec<String>,
resolved: Option<ResolvedProposal>,
}
struct ResolvedProposal {
action: ActionKind,
proposal: Proposal,
rendered: String,
summary_key: &'static str,
summary_args: serde_json::Map<String, Value>,
rollbackable: bool,
evalset_hash: Option<String>,
importance: f64,
fact_fields: Option<serde_json::Map<String, Value>>,
}
fn plan_edit_allowed(path: &str) -> bool {
let seg: Vec<&str> = path.split('.').collect();
match seg.as_slice() {
["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
["retries", node] => !node.is_empty(),
_ => false,
}
}
fn plan_get(body: &Value, path: &str) -> Value {
let mut cur = body;
for seg in path.split('.') {
cur = match cur {
Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
Some(v) => v,
None => return Value::Null,
},
Value::Object(o) => match o.get(seg) {
Some(v) => v,
None => return Value::Null,
},
_ => return Value::Null,
};
}
cur.clone()
}
fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
let segs: Vec<&str> = path.split('.').collect();
let Some((last, parents)) = segs.split_last() else {
return false;
};
let mut cur = body;
for seg in parents {
cur = match cur {
Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
Some(v) => v,
None => return false,
},
Value::Object(o) => match o.get_mut(*seg) {
Some(v) => v,
None => return false,
},
_ => return false,
};
}
match cur {
Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
Some(slot) => {
*slot = to;
true
}
None => false,
},
Value::Object(o) => {
o.insert((*last).to_string(), to);
true
}
_ => false,
}
}
fn plan_value_ok(path: &str, to: &Value) -> bool {
let seg: Vec<&str> = path.split('.').collect();
match seg.as_slice() {
["edges", _, "cond"] => to.as_str().is_some_and(|c| {
!c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
}),
["edges", _, "max_cycles"] | ["retries", _] => {
to.as_u64().is_some_and(|n| n <= 1_000)
}
_ => false,
}
}
fn resolve_proposal<S: OmsSubstrate>(
sub: &S,
d: &crate::llm::LlmDraft,
target: &TargetRef,
cited: &[String],
ns_by_hash: &std::collections::BTreeMap<String, String>,
caps: Capabilities,
policy: &crate::policy::Policy,
) -> Option<ResolvedProposal> {
use crate::llm::DraftProposal as P;
let (skills, plans) = (&policy.skills, &policy.plans);
let mut args = serde_json::Map::new();
match d.parsed_proposal()? {
P::Plan { description, when_to_use, nodes, edges } => {
if !plans.enabled || !caps.plans {
return None;
}
let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
args.insert("name".into(), Value::from(name.clone()));
args.insert("nodes".into(), Value::from(n_nodes as u64));
let mut stmts = vec![match &existing_skill {
Some(h) => cal::supersede(h, "skill", &skill),
None => cal::add("skill", &skill),
}];
let (summary_key, kind) = match &workflow {
Some(wf) => {
stmts.push(match &existing_plan {
Some(h) => cal::supersede(h, "workflow", wf),
None => cal::add("workflow", wf),
});
args.insert("edges".into(), Value::from(n_edges as u64));
("llm.plan", "plan")
}
None => {
args.insert("steps".into(), Value::from(n_nodes as u64));
("llm.skill", "skill")
}
};
let patched = existing_skill.is_some() || (workflow.is_some() && existing_plan.is_some());
let (action, verb) = if patched {
(ActionKind::Revise, "revise")
} else {
(ActionKind::Record, "record")
};
Some(ResolvedProposal {
action,
proposal: Proposal::Cal { cal: cal::batch(&stmts) },
rendered: format!(
"Proposed {kind} to {verb}: \"{name}\" — {n_nodes} steps; when: {}",
skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
),
summary_key,
summary_args: args,
rollbackable: true,
evalset_hash: None,
importance: 0.65,
fact_fields: None,
})
}
P::Skill { description, when_to_use, steps } => {
if !skills.enabled {
return None;
}
let SkillFields { fields, name, n_steps, existing } =
derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
args.insert("name".into(), Value::from(name.clone()));
args.insert("steps".into(), Value::from(n_steps as u64));
let (action, cal, verb) = match existing {
Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
None => (ActionKind::Record, cal::add("skill", &fields), "record"),
};
Some(ResolvedProposal {
action,
proposal: Proposal::Cal { cal },
rendered: format!(
"Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
),
summary_key: "llm.skill",
summary_args: args,
rollbackable: true,
evalset_hash: None,
importance: 0.6,
fact_fields: None,
})
}
P::Lesson { lesson } => {
let lesson = sanitize_lesson(&lesson);
if lesson.is_empty() {
return None;
}
let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
args.insert("lesson".into(), Value::from(lesson.clone()));
Some(ResolvedProposal {
action: ActionKind::ClusterFailure,
proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
rendered: format!("Proposed lesson to record: \"{lesson}\""),
summary_key: "llm.lesson",
summary_args: args,
rollbackable: true,
evalset_hash: None,
importance: 0.5,
fact_fields: Some(fields),
})
}
P::Fact { relation, object } => {
let relation = sanitize_relation(&relation)?;
let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
if object.is_empty() {
return None;
}
let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
args.insert("relation".into(), Value::from(relation.clone()));
args.insert("object".into(), Value::from(object.clone()));
Some(ResolvedProposal {
action: ActionKind::Record,
proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
summary_key: "llm.fact",
summary_args: args,
rollbackable: true,
evalset_hash: None,
importance: 0.5,
fact_fields: Some(fields),
})
}
P::QueryRevision { body } => {
let name = target.opaque();
if name.is_empty()
|| name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
{
return None;
}
let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
if body.is_empty() || !safe_definition_body(&body) {
return None;
}
let stmt = match target.scheme() {
"query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
"template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
_ => return None,
};
sub.validate_cal(&stmt).ok()?;
sub.definition_inverse(&stmt).ok().flatten()?;
args.insert("name".into(), Value::from(name));
args.insert("body".into(), Value::from(body.clone()));
Some(ResolvedProposal {
action: ActionKind::Revise,
proposal: Proposal::Cal { cal: stmt },
rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
summary_key: "llm.query_revision",
summary_args: args,
rollbackable: true,
evalset_hash: None,
importance: 0.6,
fact_fields: None,
})
}
P::PlanRevision { edits } => {
if !caps.plans
|| target.scheme() != "grain"
|| edits.is_empty()
|| edits.len() > crate::llm::MAX_PLAN_EDITS
{
return None;
}
let hash = target.opaque();
let g = sub.grain(hash).ok().flatten()?;
if g.grain_type != "workflow" || !g.is_live() {
return None;
}
let mut body = Value::Object(g.fields.clone());
let mut deltas = Vec::new();
let nodes: std::collections::BTreeSet<String> = body
.get("nodes")
.and_then(Value::as_array)
.map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
.unwrap_or_default();
for e in &edits {
if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
return None;
}
if let Some(node) = e.path.strip_prefix("retries.") {
if !nodes.contains(node) {
return None;
}
}
if plan_get(&body, &e.path) != e.from {
return None;
}
if e.from == e.to {
return None;
}
if !plan_set(&mut body, &e.path, e.to.clone()) {
return None;
}
deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
}
sub.validate_plan(&body).ok()?;
let Value::Object(fields) = body else {
return None;
};
let stmt = cal::supersede(hash, "workflow", &fields);
sub.validate_cal(&stmt).ok()?;
args.insert("plan".into(), Value::from(hash));
args.insert("edits".into(), Value::from(deltas.join("; ")));
Some(ResolvedProposal {
action: ActionKind::Revise,
proposal: Proposal::Cal { cal: stmt },
rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
summary_key: "llm.plan_revision",
summary_args: args,
rollbackable: true,
evalset_hash: None,
importance: 0.7,
fact_fields: None,
})
}
P::CodeRevision { source } => {
if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
return None;
}
if source.chars().count() > crate::llm::MAX_CODE_LEN {
return None;
}
let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
let mut data = serde_json::Map::new();
data.insert("tool".into(), Value::from(target.opaque()));
data.insert("source".into(), Value::from(source.clone()));
args.insert("tool".into(), Value::from(target.opaque()));
args.insert("bytes".into(), Value::from(source.len() as u64));
Some(ResolvedProposal {
action: ActionKind::CodeRevision,
proposal: Proposal::Data { data },
rendered: format!(
"Proposed new source for tool {} ({} bytes), gated by evalset {}",
target.opaque(),
source.len(),
evalset
),
summary_key: "llm.code_revision",
summary_args: args,
rollbackable: true,
evalset_hash: Some(evalset),
importance: 0.8,
fact_fields: None,
})
}
}
}
fn stamp_llm(
model: &str,
d: &crate::llm::LlmDraft,
target_ref: String,
cited: Vec<String>,
resolved: Option<ResolvedProposal>,
confidence: f64,
now_ms: i64,
) -> Recommendation {
let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
let guidance = if d.guidance.trim().is_empty() {
None
} else {
Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
};
let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
Some(mut r) => {
let content = match &r.fact_fields {
Some(fields) => format!(
"{} {}",
fields.get("relation").and_then(Value::as_str).unwrap_or(""),
fields.get("object").and_then(Value::as_str).unwrap_or("")
),
None => match &r.proposal {
Proposal::Cal { cal } => cal.clone(),
Proposal::Data { data } => Value::Object(data.clone()).to_string(),
Proposal::Edit { diff, .. } => diff.clone(),
},
};
if let Some(mut fields) = r.fact_fields.take() {
fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
r.proposal = Proposal::Cal { cal: cal::add("fact", &fields) };
}
let mut args = r.summary_args;
args.insert("text".into(), Value::from(summary_text));
(
r.action,
r.proposal,
Summary::new(r.summary_key, args),
r.rollbackable,
r.importance,
r.evalset_hash,
Some(content),
)
}
None => {
let mut args = serde_json::Map::new();
args.insert("text".into(), Value::from(summary_text));
let mut data = serde_json::Map::new();
data.insert("source".into(), Value::from("llm"));
(
ActionKind::Flag,
Proposal::Data { data },
Summary::new("llm.discover", args),
false,
0.3,
None,
None,
)
}
};
let dedup = match &content {
Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
None => dedup_key("llm", &target_ref, action),
};
Recommendation {
hash: String::new(),
analyzer: "loop.llm/1".to_string(),
params_snapshot: serde_json::Map::new(),
origin: Origin::Llm { model: model.to_string() },
target_ref: target_ref.clone(),
action_kind: action,
dedup_key: dedup,
summary,
severity: Severity::Low,
proposal,
destructive: false,
rollbackable,
evidence: cited,
evidence_query: None,
metric: None,
confidence: confidence.clamp(0.0, 1.0),
importance,
created_at_ms: now_ms,
guidance,
evalset_hash,
status: RecStatus::Pending,
}
}
fn skill_instructions(min_steps: u32) -> String {
format!(
" (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
calls that reached its goal, which a later session facing the same situation \
should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
the tool called and the values that mattered (the field checked, the tag set, the \
exact format produced), a one-line description, and 'when_to_use' — the situation \
that should trigger it. The skill-name is a short identifier (letters, digits, \
_ -). If a saved skill already covers this procedure, use ITS name so it is \
patched rather than duplicated. Do not propose a skill for a procedure that \
failed, or for one already saved and unchanged. A finding that itself describes \
two or more steps the agent should carry out in order ('after listing the \
tickets, fetch each, then …') IS a procedure: propose it as a skill or a plan, \
never as a lesson — a lesson is one rule, and a procedure written as one is a \
procedure nobody can open.",
crate::llm::MAX_SKILL_STEPS
)
}
fn plan_instructions(min_nodes: u32) -> String {
format!(
" (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
the cited evidence — and what the step does with it; and 'edges' from step to \
step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
'path != literal', 'path exists' or '!path', where path is dotted names and the \
literal is a JSON string, number, true, false or null — no other operators; state \
a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
has branches or a loop; prefer a skill when it is a straight list. If a saved \
plan already covers this procedure, use ITS name so it is patched.",
crate::llm::MAX_PLAN_NODES
)
}
struct PlanFields {
skill: serde_json::Map<String, Value>,
workflow: Option<serde_json::Map<String, Value>>,
name: String,
n_nodes: usize,
n_edges: usize,
existing_skill: Option<String>,
existing_plan: Option<String>,
}
#[allow(clippy::too_many_arguments)]
fn derived_plan_fields<S: SubstrateRead>(
sub: &S,
target: &TargetRef,
description: &str,
when_to_use: &str,
nodes: &[crate::llm::PlanNodeDraft],
edges: &[crate::llm::PlanEdgeDraft],
cited: &[String],
ns_by_hash: &std::collections::BTreeMap<String, String>,
plans: &crate::policy::PlanAuthoring,
) -> Option<PlanFields> {
if target.scheme() != "entity" {
return None;
}
let name = sanitize_skill_name(
target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
)?;
let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
if description.is_empty() || when_to_use.is_empty() {
return None;
}
if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
return None;
}
let known_tools: BTreeSet<String> = cited
.iter()
.filter_map(|h| sub.grain(h).ok().flatten())
.filter_map(|g| g.tool_name().map(normalize_ident))
.collect();
let mut ids: Vec<String> = Vec::new();
let mut steps: Vec<String> = Vec::new();
let mut seen: BTreeSet<String> = BTreeSet::new();
let mut grounded = 0usize;
for n in nodes {
let id = sanitize_skill_name(&n.id)?;
if !seen.insert(id.clone()) {
return None; }
let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
if step.is_empty() {
return None;
}
let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
if !tool.is_empty() && known_tools.contains(&normalize_ident(&tool)) {
grounded += 1;
steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
} else {
steps.push(format!("{}. {id}: {step}", steps.len() + 1));
}
ids.push(id);
}
if grounded == 0 {
return None;
}
let mut edge_vals: Vec<Value> = Vec::new();
let mut flow_lines: Vec<String> = Vec::new();
let mut runnable = true;
for e in edges {
let (Some(src), Some(dst)) = (sanitize_skill_name(&e.src), sanitize_skill_name(&e.dst)) else {
runnable = false;
continue;
};
if !seen.contains(&src) || !seen.contains(&dst) {
flow_lines.push(format!("{} → {}", e.src.trim(), e.dst.trim()));
runnable = false;
continue;
}
let mut ev = serde_json::Map::new();
ev.insert("src".into(), Value::from(src.clone()));
ev.insert("dst".into(), Value::from(dst.clone()));
let mut label = format!("{src} → {dst}");
if let Some(c) = e
.cond
.as_deref()
.map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
.filter(|c| !c.is_empty())
{
label.push_str(&format!(" if {c}"));
ev.insert("cond".into(), Value::from(c));
}
if let Some(m) = e.max_cycles {
if m == 0 || m > 100 {
runnable = false;
} else {
label.push_str(&format!(" (at most {m} times)"));
ev.insert("max_cycles".into(), Value::from(m));
}
}
flow_lines.push(label);
edge_vals.push(Value::Object(ev));
}
if edge_vals.len() > 4 * ids.len() {
runnable = false;
}
let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
for h in cited {
if let Some(ns) = ns_by_hash.get(h) {
if !ns.is_empty() {
*ns_counts.entry(ns.as_str()).or_default() += 1;
}
}
}
let ns = ns_counts
.iter()
.max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
.map(|(ns, _)| ns.to_string());
let mut workflow = serde_json::Map::new();
workflow.insert("nodes".into(), Value::from(ids.clone()));
workflow.insert("edges".into(), Value::Array(edge_vals));
workflow.insert("name".into(), Value::from(name.clone()));
if let Some(ns) = &ns {
workflow.insert("namespace".into(), Value::from(ns.clone()));
}
let workflow = (runnable && sub.validate_plan(&Value::Object(workflow.clone())).is_ok())
.then_some(workflow);
let mut instructions = steps.join("\n");
if !flow_lines.is_empty() {
instructions.push_str("\n\nFlow:\n");
instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
}
let mut skill = serde_json::Map::new();
skill.insert("name".into(), Value::from(name.clone()));
skill.insert("description".into(), Value::from(description));
skill.insert("when_to_use".into(), Value::from(when_to_use));
skill.insert("instructions".into(), Value::from(instructions));
if let Some(ns) = &ns {
skill.insert("namespace".into(), Value::from(ns.clone()));
}
let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
.ok()?
.into_iter()
.find(|g| pick(g))
.map(|g| g.hash)
};
let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
let existing_plan = workflow
.is_some()
.then(|| live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str())))
.flatten();
let n_edges = workflow
.as_ref()
.and_then(|w| w.get("edges"))
.and_then(Value::as_array)
.map_or(0, |a| a.len());
Some(PlanFields { skill, workflow, name, n_nodes: ids.len(), n_edges, existing_skill, existing_plan })
}
pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
fn detect_premise_drift<S: OmsSubstrate>(
sub: &S,
p: &mut LoopPersisted,
now_ms: i64,
) -> Result<Vec<OutcomeInput>> {
let mut out = Vec::new();
let applied: Vec<(String, String, Vec<String>)> = p
.applied
.iter()
.filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
.map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
.collect();
for (rec_hash, target_ref, own) in applied {
let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
if rec.evidence.is_empty() {
continue;
}
let mut moved = 0u64;
for e in &rec.evidence {
match sub.grain(e)? {
None => moved += 1, Some(g) => {
let Some(newer) = &g.superseded_by else { continue };
if own.iter().any(|c| c == newer) {
continue; }
match sub.grain(newer)? {
None => moved += 1,
Some(n) => {
if !same_value(&g, &n) {
moved += 1;
}
}
}
}
}
}
if moved == 0 {
continue;
}
let already = p
.outcomes
.get(&rec_hash)
.and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
.is_some_and(|o| o.current == moved as f64);
if !already {
p.outcomes.entry(rec_hash.clone()).or_default().push(
crate::recommendation::OutcomeResult {
rec_hash: rec_hash.clone(),
metric: PREMISE_DRIFT_METRIC.into(),
baseline: 0.0,
current: moved as f64,
verdict: "drifted".into(),
horizon_ms: 0,
checkpoint: None,
measured_at_ms: now_ms,
},
);
}
out.push(OutcomeInput {
rec_hash,
target_ref,
metric: PREMISE_DRIFT_METRIC.into(),
baseline: 0.0,
current: moved as f64,
unit: "superseded premises".into(),
higher_is_better: false,
});
}
Ok(out)
}
fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
return normalize_ident(a) == normalize_ident(b);
}
for key in ["content", "tool_content", "body", "text", "object"] {
if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
return normalize_ident(a) == normalize_ident(b);
}
}
false
}
fn sanitize_skill_name(s: &str) -> Option<String> {
let t = s.trim();
if t.is_empty()
|| t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
|| !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
{
return None;
}
Some(t.to_string())
}
struct SkillFields {
fields: serde_json::Map<String, Value>,
name: String,
n_steps: usize,
existing: Option<String>,
}
#[allow(clippy::too_many_arguments)]
fn derived_skill_fields<S: SubstrateRead>(
sub: &S,
target: &TargetRef,
description: &str,
when_to_use: &str,
steps: &[String],
cited: &[String],
ns_by_hash: &std::collections::BTreeMap<String, String>,
skills: &crate::policy::SkillAuthoring,
) -> Option<SkillFields> {
if target.scheme() != "entity" {
return None;
}
let name = sanitize_skill_name(
target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
)?;
let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
let steps: Vec<String> = steps
.iter()
.map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
.filter(|st| !st.is_empty())
.take(crate::llm::MAX_SKILL_STEPS)
.collect();
if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
return None;
}
let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
for h in cited {
if let Some(ns) = ns_by_hash.get(h) {
if !ns.is_empty() {
*ns_counts.entry(ns.as_str()).or_default() += 1;
}
}
}
let ns = ns_counts
.iter()
.max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
.map(|(ns, _)| ns.to_string());
let instructions = steps
.iter()
.enumerate()
.map(|(i, st)| format!("{}. {st}", i + 1))
.collect::<Vec<_>>()
.join("\n");
let mut fields = serde_json::Map::new();
fields.insert("name".into(), Value::from(name.clone()));
fields.insert("description".into(), Value::from(description));
fields.insert("when_to_use".into(), Value::from(when_to_use));
fields.insert("instructions".into(), Value::from(instructions));
if let Some(ns) = &ns {
fields.insert("namespace".into(), Value::from(ns.clone()));
}
let existing = sub
.grains_of_type(
crate::model::grain_type::SKILL,
ns.as_deref(),
ReadOpts { live_only: true, since_ms: None },
)
.ok()?
.into_iter()
.find(|g| g.skill_name() == Some(name.as_str()))
.map(|g| g.hash);
Some(SkillFields { fields, name, n_steps: steps.len(), existing })
}
fn derived_fact_fields(
target: &TargetRef,
relation: &str,
object: &str,
cited: &[String],
ns_by_hash: &std::collections::BTreeMap<String, String>,
) -> Option<serde_json::Map<String, Value>> {
if target.scheme() != "entity" {
return None;
}
let subject = target
.opaque()
.rsplit_once('/')
.map(|(_, s)| s)
.unwrap_or(target.opaque());
if subject.is_empty() {
return None;
}
let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
for h in cited {
if let Some(ns) = ns_by_hash.get(h) {
if !ns.is_empty() {
*ns_counts.entry(ns.as_str()).or_default() += 1;
}
}
}
let lesson_ns = ns_counts
.iter()
.max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
.map(|(ns, _)| ns.to_string());
let mut fields = serde_json::Map::new();
fields.insert("subject".into(), Value::from(subject));
fields.insert("relation".into(), Value::from(relation));
fields.insert("object".into(), Value::from(object));
if let Some(ns) = lesson_ns {
fields.insert("namespace".into(), Value::from(ns));
}
Some(fields)
}
fn requires_gating(kind: ActionKind) -> bool {
matches!(
kind,
ActionKind::CodeRevision | ActionKind::AdapterRevision
)
}
fn stamp(
m: &AnalyzerManifest,
params: &crate::manifest::Params,
d: crate::recommendation::RecDraft,
now_ms: i64,
) -> Result<Recommendation> {
let target = TargetRef::parse(&d.target_ref)?;
crate::recommendation::validate_code_rules(
d.action_kind,
target.target_class(),
d.evalset_hash.as_deref(),
)?;
let revert_of = match (&d.action_kind, &d.proposal) {
(ActionKind::Revert, Proposal::Data { data }) => {
data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
}
_ => None,
};
let dedup = match revert_of.as_deref() {
Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
None => dedup_key(m.family(), &d.target_ref, d.action_kind),
};
let destructive = match &d.proposal {
Proposal::Cal { cal } => cal::contains_destructive(cal),
_ => false,
};
let rollbackable = match &d.proposal {
Proposal::Cal { .. } => !destructive,
Proposal::Edit { .. } => false,
Proposal::Data { .. } => requires_gating(d.action_kind),
};
let mut evidence = d.evidence;
evidence.truncate(MAX_EVIDENCE);
let origin = match m.trust_class {
crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
_ => Origin::Builtin,
};
Ok(Recommendation {
hash: String::new(),
analyzer: m.id.clone(),
params_snapshot: params.snapshot(),
origin,
target_ref: target.as_string(),
action_kind: d.action_kind,
dedup_key: dedup,
summary: d.summary,
severity: d.severity,
proposal: d.proposal,
destructive,
rollbackable,
evidence,
evidence_query: d.evidence_query,
metric: d.metric,
confidence: d.confidence,
importance: d.importance,
created_at_ms: now_ms,
guidance: None,
evalset_hash: d.evalset_hash,
status: RecStatus::Pending,
})
}
fn validate_because(because: &str) -> Result<String> {
let trimmed = because.trim();
if trimmed.is_empty() {
return Err(Error::InvalidProposal(
"a BECAUSE reason is required".into(),
));
}
if trimmed.chars().count() > MAX_BECAUSE {
return Err(Error::InvalidProposal(format!(
"BECAUSE exceeds {MAX_BECAUSE} chars"
)));
}
Ok(trimmed.to_string())
}
fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
for req in &m.requires {
match req {
Capability::Forks if !caps.forks => return Some("forks"),
Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
_ => {}
}
}
None
}
fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
p.config.get(analyzer_id).and_then(|c| c.severity_floor)
}
fn gate(
opts: &RunOptions,
p: &LoopPersisted,
new_grains: u64,
new_errors: u64,
now_ms: i64,
) -> Option<SkipReason> {
let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
if !any {
return None;
}
let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
let stale_ok = opts
.if_stale_ms
.is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
if min_new_ok || min_err_ok || stale_ok {
return None;
}
if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
Some(SkipReason::NotStale)
} else {
Some(SkipReason::MinNewNotMet)
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub(crate) struct NewSince {
pub grains: u64,
pub error_events: u64,
pub events: u64,
pub sessions: u64,
}
fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
let opts = ReadOpts {
live_only: false,
since_ms: watermark.map(|w| w + 1),
};
let mut n = NewSince::default();
let mut sessions: BTreeSet<&str> = BTreeSet::new();
let mut events_held: Vec<GrainRecord> = Vec::new();
for t in [
crate::model::grain_type::FACT,
crate::model::grain_type::EVENT,
crate::model::grain_type::TOOL,
crate::model::grain_type::OBSERVATION,
] {
let g = sub.grains_of_type(t, None, opts)?;
n.grains += g.len() as u64;
if t == crate::model::grain_type::TOOL {
n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
}
if t == crate::model::grain_type::EVENT {
n.events = g.len() as u64;
events_held = g;
}
}
for e in &events_held {
if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
sessions.insert(sid);
}
}
n.sessions = sessions.len() as u64;
Ok(n)
}
fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
if !c.is_set() {
return None;
}
let time_ok = c
.every_ms
.is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
let events_ok = c.every_events.is_some_and(|m| new.events >= m);
let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
if time_ok || grains_ok || events_ok || sessions_ok {
None
} else {
Some(SkipReason::CadenceNotDue)
}
}
fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
let grains = sub.grains_of_type(
crate::model::grain_type::RECOMMENDATION,
Some(LOOP_NS),
ReadOpts {
live_only: false,
since_ms: None,
},
)?;
let mut set = BTreeSet::new();
for g in grains {
let status = p
.status_index
.get(&g.hash)
.copied()
.unwrap_or(RecStatus::Pending);
if matches!(
status,
RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
) {
if let Some(key) = g.str_field("dedup_key") {
set.insert(key.to_string());
}
}
}
Ok(set)
}
fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
const BASE_MS: i64 = 7 * 86_400_000;
const CAP_MS: i64 = 90 * 86_400_000;
let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
*strikes = strikes.saturating_add(1);
p.cooldowns.insert(dedup_key, now_ms + interval);
}
fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
let g = sub
.grain(rec_hash)?
.ok_or_else(|| Error::NotFound(rec_hash.into()))?;
Recommendation::from_fields(rec_hash, &g.fields)
}
pub(crate) fn is_definition_statement(line: &str) -> bool {
let up = line.trim_start().to_ascii_uppercase();
up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
}
const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
approve it to acknowledge it and let it expire.";
const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
(evalset hash + run id + stats) — use apply_gated";
const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
acknowledge it and let it expire.";
pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
match proposal {
Proposal::Cal { .. } => Ok(()),
Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
Proposal::Data { data } => {
if requires_gating(action_kind)
|| data.get("revert_of").and_then(Value::as_str).is_some()
{
Ok(())
} else {
Err(Error::InvalidProposal(ADVISORY_DATA.into()))
}
}
}
}
#[cfg(test)]
mod definition_body_tests {
use super::safe_definition_body;
#[test]
fn ordinary_bodies_pass() {
assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
}
#[test]
fn a_body_cannot_close_its_own_block_or_carry_destruction() {
assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
assert!(!safe_definition_body("RECALL facts FORGET abc"));
assert!(!safe_definition_body("recall facts purge older than 1d"));
assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
}
}
#[cfg(test)]
mod plan_edit_tests {
use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
use serde_json::json;
fn plan() -> serde_json::Value {
json!({
"nodes": ["fetch", "review", "post"],
"edges": [
{"src": "fetch", "dst": "review"},
{"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
],
"bindings": {"fetch": "sha256:tool1"},
"retries": {"fetch": 1}
})
}
#[test]
fn the_allowlist_admits_thresholds_and_refuses_topology() {
assert!(plan_edit_allowed("edges.1.cond"));
assert!(plan_edit_allowed("edges.1.max_cycles"));
assert!(plan_edit_allowed("retries.fetch"));
for path in [
"nodes",
"nodes.0",
"edges.0.src",
"edges.0.dst",
"edges",
"bindings.fetch",
"edges.x.cond",
"",
] {
assert!(!plan_edit_allowed(path), "{path} must not be editable");
}
}
#[test]
fn values_are_type_checked_against_the_field() {
assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
assert!(plan_value_ok("retries.fetch", &json!(3)));
assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
assert!(!plan_value_ok("edges.0.src", &json!("other")));
}
#[test]
fn get_reads_through_arrays_and_objects_and_absence_is_null() {
let p = plan();
assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
assert_eq!(plan_get(&p, "retries.review"), json!(null));
assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
}
#[test]
fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
let mut p = plan();
assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
assert!(plan_set(&mut p, "retries.review", json!(2)));
assert_eq!(plan_get(&p, "retries.review"), json!(2));
assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
}
}
#[cfg(test)]
mod definition_proposal_tests {
use super::is_definition_statement;
#[test]
fn definition_statements_are_recognized_in_both_spellings() {
assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
assert!(is_definition_statement(" define template foo AS { x }"));
assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
assert!(!is_definition_statement("ADD fact {}"));
assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
assert!(!is_definition_statement("FORGET abc"));
assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
}
}