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,
};
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,
}
#[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>,
}
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,
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>,
}
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_grains, new_error_events) = count_new(sub, watermark)?;
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 outcome_inputs = measure_outcomes(sub, &mut persisted, now_ms)?;
let AnalysisPass {
survivors,
proposed,
deduped,
analyzers_run,
analyzers_skipped,
} = 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,
})
}
#[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(),
}),
}
}
if self.llm.is_some() {
candidates.extend(self.discover(
sub,
&candidates,
analysis_watermark,
&opts.namespaces,
now_ms,
));
}
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,
})
}
fn discover<S: OmsSubstrate>(
&self,
sub: &S,
candidates: &[Recommendation],
watermark: Option<i64>,
namespaces: &[String],
now_ms: i64,
) -> 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 mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
let mut bundle: BTreeSet<String> = BTreeSet::new();
for c in candidates {
for h in &c.evidence {
if !bundle.contains(h) {
if let Ok(Some(g)) = sub.grain(h) {
push_evidence(&mut evidence, &mut bundle, &g);
}
}
}
}
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 };
'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() >= 64 {
break 'seed;
}
push_evidence(&mut evidence, &mut bundle, &g);
}
}
}
}
if evidence.is_empty() {
return Vec::new(); }
let (approved, rejected) = self.llm_history(sub);
let request = crate::llm::LlmRequest {
loop_proto: 1,
op: "discover",
instructions: DISCOVER_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 mut validated: Vec<(crate::llm::LlmDraft, String, Vec<String>)> = Vec::new();
for d in crate::llm::parse_discover(&raw)
.recommendations
.into_iter()
.take(crate::llm::MAX_LLM_DRAFTS)
{
let cited: Vec<String> =
d.evidence.iter().filter(|h| bundle.contains(*h)).cloned().collect();
if cited.is_empty() {
continue; }
let Ok(target) = TargetRef::parse(&d.target) else {
continue;
};
let tc = target.target_class();
if tc != "memory" && tc != "query" {
continue; }
validated.push((d, target.as_string(), cited));
}
if validated.is_empty() {
return Vec::new();
}
let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
self.verify_drafts(&**llm, ground, &validated, &evidence, now_ms)
}
fn verify_drafts(
&self,
llm: &dyn crate::llm::LlmBackend,
ground: &dyn crate::llm::LlmBackend,
validated: &[(crate::llm::LlmDraft, String, Vec<String>)],
evidence: &[crate::llm::EvidenceItem],
now_ms: i64,
) -> 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, (d, _t, cited))| GroundItem {
id: i,
claim: cap(&d.summary, MAX_SUMMARY_LEN),
evidence: ev_for(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) => parse_ground(&raw)
.results
.into_iter()
.filter(|r| r.supported)
.map(|r| r.id)
.collect(),
None => return Vec::new(),
};
if grounded.is_empty() {
return Vec::new();
}
let items: Vec<VerifyItem> = validated
.iter()
.enumerate()
.filter(|(i, _)| grounded.contains(i))
.map(|(i, (d, t, cited))| VerifyItem {
id: i,
summary: cap(&d.summary, MAX_SUMMARY_LEN),
target: t.clone(),
evidence: ev_for(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(),
};
let mut out = Vec::new();
for (i, (d, target_str, cited)) in validated.iter().enumerate() {
if let Some(&conf) = verdicts.get(&i) {
if conf >= MIN_LLM_CONFIDENCE {
out.push(stamp_llm(
llm.model(),
d,
target_str.clone(),
cited.clone(),
conf,
now_ms,
));
}
}
}
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) {
const BASE_MS: i64 = 7 * 86_400_000;
const CAP_MS: i64 = 90 * 86_400_000;
let strikes = p.cooldown_strikes.entry(rec.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(rec.dedup_key, now_ms + interval);
}
}
sub.store_state(&p.to_value()?)?;
Ok(())
}
pub fn preflight_apply<S: OmsSubstrate>(
&self,
sub: &S,
rec_hash: &str,
scopes: &ScopeSet,
allow_destructive: 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.proposal)?;
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,
)
}
#[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 rec.action_kind == ActionKind::CodeRevision {
let g = gating.ok_or_else(|| {
Error::InvalidProposal(
"code revisions apply only with a recorded gating run \
(evalset hash + run id + stats) — use apply_gated"
.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 \
cannot admit code",
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 rec.action_kind == ActionKind::CodeRevision => {
let mut spec = crate::substrate::GrainSpec::new(
crate::model::grain_type::FACT,
LOOP_NS,
)
.with_field("subject", rec.target_ref.clone())
.with_field("relation", "mg:code_promotion")
.with_field(
"object",
serde_json::to_string(data)
.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()?)?;
}
}
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 (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
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, i64)> = 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 horizon in metric.horizons() {
if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
due.push((h.clone(), a.clone(), horizon));
}
}
}
let mut out = Vec::new();
for (rec_hash, applied, horizon) in due {
let metric = applied.metric.as_ref().unwrap();
let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
continue; };
let regressed = crate::recommendation::is_regression(
metric.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: metric.baseline,
current,
verdict: if regressed { "regressed" } else { "held" }.into(),
horizon_ms: horizon,
measured_at_ms: now_ms,
},
);
p.measured.entry(rec_hash.clone()).or_default().push(horizon);
if regressed {
out.push(OutcomeInput {
rec_hash,
target_ref: applied.target_ref.clone(),
metric: metric.metric.clone(),
baseline: metric.baseline,
current,
unit: metric.unit.clone(),
higher_is_better: metric.higher_is_better,
});
}
}
Ok(out)
}
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(match field {
"failed" => Some(run.failed as f64),
"passed" => Some(run.passed as f64),
"total" => Some(run.total() as f64),
"error_rate" => match run.total() {
0 => None, t => Some(run.failed as f64 / t as f64),
},
other => run.field(other),
})
}
_ => 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 MIN_LLM_CONFIDENCE: f64 = 0.75;
const DISCOVER_INSTRUCTIONS: &str = "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). 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. 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 hashes from the bundle, target a memory entity, and include \
your confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
\"target\":\"entity:<ns>/<subject>\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
\"confidence\":0.0}]}. Propose nothing you cannot ground in the evidence.";
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>,
g: &GrainRecord,
) {
if evidence.len() < 64 && bundle.insert(g.hash.clone()) {
evidence.push(crate::llm::EvidenceItem {
hash: g.hash.clone(),
grain_type: g.grain_type.clone(),
text: crate::llm::cap(&grain_brief(g), 400),
});
}
}
fn grain_brief(g: &GrainRecord) -> String {
if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
return format!("{s} {r} {o}");
}
for key in ["content", "body", "text", "summary"] {
if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
if !v.is_empty() {
return v.to_string();
}
}
}
String::new()
}
fn stamp_llm(
model: &str,
d: &crate::llm::LlmDraft,
target_ref: String,
cited: Vec<String>,
confidence: f64,
now_ms: i64,
) -> Recommendation {
let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
let mut args = serde_json::Map::new();
args.insert("text".into(), Value::from(summary_text));
let guidance = if d.guidance.trim().is_empty() {
None
} else {
Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
};
let action = ActionKind::Flag;
let mut data = serde_json::Map::new();
data.insert("source".into(), Value::from("llm"));
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_key("llm", &target_ref, action),
summary: Summary::new("llm.discover", args),
severity: Severity::Low,
proposal: Proposal::Data { data },
destructive: false,
rollbackable: false,
evidence: cited,
evidence_query: None,
metric: None,
confidence: confidence.clamp(0.0, 1.0),
importance: 0.3,
created_at_ms: now_ms,
guidance,
evalset_hash: None,
status: RecStatus::Pending,
}
}
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 dedup = 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 { .. } => d.action_kind == ActionKind::CodeRevision,
};
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)
}
}
fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
let opts = ReadOpts {
live_only: false,
since_ms: watermark.map(|w| w + 1),
};
let mut new_grains = 0u64;
let mut new_errors = 0u64;
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)?;
new_grains += g.len() as u64;
if t == crate::model::grain_type::TOOL {
new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
}
}
Ok((new_grains, new_errors))
}
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 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 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(proposal: &Proposal) -> Result<()> {
match proposal {
Proposal::Cal { .. } => Ok(()),
Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
Proposal::Data { data } => {
if data.get("revert_of").and_then(Value::as_str).is_some() {
Ok(())
} else {
Err(Error::InvalidProposal(ADVISORY_DATA.into()))
}
}
}
}
#[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""#));
}
}