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>,
#[serde(default)]
pub withdrawn: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub decider: Option<crate::decide::DeciderReport>,
}
#[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,
#[serde(default, skip_serializing_if = "is_zero")]
pub dropped_near_duplicate: 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![],
withdrawn: 0,
decider: None,
}
}
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>>,
decider: Option<crate::decide::Decider>,
}
pub(crate) struct AnalysisPass {
pub(crate) survivors: Vec<Recommendation>,
proposed: u64,
deduped: u64,
analyzers_run: Vec<String>,
pub(crate) analyzers_skipped: Vec<AnalyzerSkip>,
llm_funnel: Option<LlmFunnel>,
decider: Option<crate::decide::DeciderReport>,
}
impl Engine {
pub fn with_builtins() -> Self {
Engine {
analyzers: crate::analyzer::builtin_analyzers(),
policy: crate::policy::Policy::default(),
llm: None,
ground_llm: None,
decider: None,
}
}
pub fn empty() -> Self {
Engine {
analyzers: vec![],
policy: crate::policy::Policy::default(),
llm: None,
ground_llm: None,
decider: 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 with_decider(mut self, backend: Box<dyn crate::decide::DecideBackend>) -> Self {
let cap = self.decider.as_ref().map(|d| d.pair_cap());
let mut d = crate::decide::Decider::new(backend);
if let Some(cap) = cap {
d = d.with_pair_cap(cap);
}
self.decider = Some(d);
self
}
pub fn with_decider_pair_cap(mut self, cap: usize) -> Self {
self.decider = self.decider.take().map(|d| d.with_pair_cap(cap));
self
}
pub fn decider(&self) -> Option<&crate::decide::Decider> {
self.decider.as_ref()
}
pub fn policy(&self) -> &crate::policy::Policy {
&self.policy
}
pub fn has_llm(&self) -> bool {
self.llm.is_some()
}
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, &self.policy, now_ms)?;
let mut withdrawn = 0u64;
if self.policy.premise_drift {
outcome_inputs.extend(detect_premise_drift(sub, &mut persisted, now_ms)?);
withdrawn = withdraw_drifted_open(
sub,
&mut persisted,
self.policy.premise_drift_open_all,
now_ms,
)?;
}
let AnalysisPass {
survivors,
proposed,
deduped,
analyzers_run,
analyzers_skipped,
llm_funnel,
decider,
} = 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,
withdrawn,
decider,
})
}
#[allow(clippy::too_many_arguments)]
#[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)?;
self.analysis_pass_inner(
sub,
persisted,
&self.policy,
opts,
external_overrides,
analysis_watermark,
now_ms,
outcome_inputs,
&existing,
None,
)
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn analysis_pass_inner<S: OmsSubstrate>(
&self,
sub: &S,
persisted: &LoopPersisted,
policy: &crate::policy::Policy,
opts: &RunOptions,
external_overrides: &BTreeMap<String, Map<String, Value>>,
analysis_watermark: Option<i64>,
now_ms: i64,
outcome_inputs: &[OutcomeInput],
existing: &BTreeSet<String>,
replay: Option<&str>,
) -> Result<AnalysisPass> {
let mut analyzers_run = Vec::new();
let mut analyzers_skipped = Vec::new();
let mut candidates: Vec<Recommendation> = Vec::new();
let caps = sub.capabilities();
let verdicts = latest_verdicts(persisted);
let decider = if replay.is_none() { self.decider.as_ref() } else { None };
if let Some(d) = decider {
d.reset();
}
for analyzer in &self.analyzers {
let m = analyzer.manifest();
if let Some(why) = replay {
let out_of_process = m.trust_class == crate::manifest::TrustClass::Command;
let telemetry_fed = m.requires.contains(&crate::manifest::Capability::Telemetry);
if out_of_process || telemetry_fed {
analyzers_skipped.push(AnalyzerSkip {
id: m.id.clone(),
reason: format!(
"not replayed: {}",
if out_of_process { why } else { "telemetry rollups are not time-indexed" }
),
});
continue;
}
}
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 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,
&verdicts,
)
.with_decider(decider);
match analyzer.analyze(&ctx) {
Ok(drafts) => {
analyzers_run.push(m.id.clone());
for draft in drafts {
match stamp(m, ¶ms, draft, now_ms, ns_slice) {
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() && replay.is_none() {
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),
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() && replay.is_none() {
self.enrich(&mut survivors);
}
Ok(AnalysisPass {
survivors,
proposed,
deduped,
analyzers_run,
analyzers_skipped,
llm_funnel: self.llm.is_some().then_some(funnel),
decider: decider.map(|d| d.report()),
})
}
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));
}
if findings.iter().any(|f| f.analyzer.starts_with("loop.lesson_pile/")) {
instructions.push_str(CONSOLIDATION_INSTRUCTIONS);
}
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(
sub, &**llm, ground, validated, &evidence, outcome_metric, now_ms, funnel, namespaces,
)
}
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)]
#[allow(clippy::too_many_arguments)]
fn verify_drafts<S: SubstrateRead>(
&self,
sub: &S,
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,
scope: &[String],
) -> 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 decided = self.decide_ground_verify(&validated, &claims);
let ground_req = GroundRequest {
loop_proto: 1,
op: "ground",
instructions: GROUND_INSTRUCTIONS,
claims,
};
let grounded: std::collections::BTreeSet<usize> = if let Some(d) = &decided {
funnel.ground_verdicts = d.len() as u64;
d.iter().filter(|(_, j)| j.grounded).map(|(i, _)| *i).collect()
} else {
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(&self_report) = verdicts.get(&i) {
let judged = decided.as_ref().and_then(|d| d.get(&i));
let conf = judged.map_or(self_report, |j| j.sound);
if conf >= MIN_LLM_CONFIDENCE {
let near = match v.resolved.as_ref() {
Some(r) if r.action != ActionKind::Consolidate => r
.fact_fields
.as_ref()
.filter(|f| f.get("relation").and_then(Value::as_str) == Some("lesson"))
.map(|f| {
near_duplicates_of(
sub,
f.get("subject").and_then(Value::as_str).unwrap_or(""),
f.get("namespace").and_then(Value::as_str),
f.get("object").and_then(Value::as_str).unwrap_or(""),
)
})
.unwrap_or_default(),
_ => Vec::new(),
};
if !near.is_empty()
&& self.policy.near_duplicate == crate::policy::NearDuplicateMode::Suppress
{
funnel.dropped_near_duplicate += 1;
continue;
}
let mut rec = stamp_llm(
llm.model(),
&v.draft,
v.target_ref,
v.cited,
v.resolved,
conf,
now_ms,
scope,
);
if let Some(j) = judged {
rec.llm_confidence = Some(self_report);
rec.judged_by = Some(j.judged_by.clone());
}
if let Some(best) = near.first() {
rec.summary.args.insert("near_count".into(), Value::from(near.len() as u64));
rec.summary.args.insert("near_score".into(), Value::from(best.score));
rec.summary.args.insert("near_method".into(), Value::from(best.method.clone()));
rec.summary.args.insert(
"near_hash".into(),
Value::from(best.hash.chars().take(12).collect::<String>()),
);
rec.summary.template_id = "llm.lesson_near_duplicate".into();
rec.near_duplicate_of = near;
}
if rec.rollbackable {
rec.metric = outcome_metric.clone();
}
out.push(rec);
}
}
}
funnel.stored = out.len() as u64;
out
}
fn decide_ground_verify(
&self,
validated: &[ValidatedDraft],
claims: &[crate::llm::GroundItem],
) -> Option<BTreeMap<usize, DecidedDraft>> {
use crate::decide::{Ask, DECIDE_MIN_P};
let d = self.decider.as_ref()?;
if !d.calibrated() {
return None;
}
let backend = d.describe();
let mut out = BTreeMap::new();
for (i, (v, c)) in validated.iter().zip(claims).enumerate() {
let evidence: Vec<Value> = c
.evidence
.iter()
.map(|e| serde_json::json!({"id": e.id, "grain_type": e.grain_type, "text": e.text}))
.collect();
let state = serde_json::json!({
"recommendation": {"summary": c.claim, "guidance": v.draft.guidance},
"evidence": evidence,
});
let mut asks: Vec<Ask> = c
.evidence
.iter()
.map(|e| Ask::Noul {
id: format!("ev_{}", e.id),
instructions: format!(
"Does evidence item \"{}\" (in state.evidence) contain the premise the recommendation relies on?",
e.id
),
})
.collect();
asks.push(Ask::Noul {
id: "sound".into(),
instructions: "Given only this evidence, is the recommendation sound?".into(),
});
let a = d.ask(state, &asks).ok().filter(|a| a.calibrated)?;
let sound = *a.noul.get("sound")?;
let grounded = a
.noul
.iter()
.any(|(id, p)| id.starts_with("ev_") && *p >= DECIDE_MIN_P);
let judged_by = a.judged_by(&backend, "ground_verify", a.noul.clone());
out.insert(i, DecidedDraft { grounded, sound, judged_by });
}
Some(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;
}
if rec.judged_by.is_some() {
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 | RecStatus::Withdrawn => {}
}
}
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,
policy: &crate::policy::Policy,
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 bound = cost_bound_for(policy, metric);
let base = baseline_at_apply(
sub,
metric,
applied.applied_at_ms,
baseline_kind_for(policy, metric),
bound.map(|b| b.field.as_str()),
)?;
let baseline = base.value;
let tolerance = tolerance_for(policy, metric, &base);
let regressed = crate::recommendation::is_regression(
baseline,
current,
metric.higher_is_better,
tolerance,
);
let (current_run_id, cost) = current_run_and_cost(sub, metric, applied.applied_at_ms, bound, &base)?;
let costlier = cost.as_ref().is_some_and(|c| c.breached());
let verdict = if regressed {
"regressed"
} else if costlier {
"held_costlier"
} else {
"held"
};
p.outcomes.entry(rec_hash.clone()).or_default().push(
crate::recommendation::OutcomeResult {
rec_hash: rec_hash.clone(),
metric: metric.metric.clone(),
baseline,
current,
verdict: verdict.into(),
baseline_kind: base.kind.into(),
baseline_run_id: base.run_id.clone(),
best_before: base.best_before,
tolerance,
current_run_id: current_run_id.clone(),
cost: cost.clone(),
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 || costlier {
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,
baseline_kind: base.kind.into(),
baseline_run_id: base.run_id,
best_before: base.best_before,
tolerance,
current_run_id,
cost,
});
}
}
Ok(out)
}
fn cost_bound_for<'p>(
policy: &'p crate::policy::Policy,
metric: &crate::recommendation::MetricSnapshot,
) -> Option<&'p crate::policy::CostBound> {
match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
(Some(e), Some((hash, _))) if e.hash == hash => e.cost.as_ref(),
_ => None,
}
}
fn current_run_and_cost<S: SubstrateRead>(
sub: &S,
metric: &crate::recommendation::MetricSnapshot,
applied_at_ms: i64,
bound: Option<&crate::policy::CostBound>,
base: &BaselineRead,
) -> Result<(Option<String>, Option<crate::recommendation::CostRead>)> {
let Some((evalset, _)) = crate::eval::parse_evalset_metric(&metric.metric) else {
return Ok((None, None));
};
let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(applied_at_ms))? else {
return Ok((None, None));
};
let cost = bound.map(|b| {
let current = crate::eval::run_value(&run, &b.field);
let status = match (base.cost, current) {
(Some(bl), Some(cur)) if cur > bl * b.max_increase_ratio + 1e-9 => "breached",
(Some(_), Some(_)) => "within",
_ => "not_measurable",
};
crate::recommendation::CostRead {
field: b.field.clone(),
max_increase_ratio: b.max_increase_ratio,
baseline: base.cost,
current,
status: status.into(),
}
});
Ok((Some(run.run_id), cost))
}
fn tolerance_for(
policy: &crate::policy::Policy,
metric: &crate::recommendation::MetricSnapshot,
base: &BaselineRead,
) -> f64 {
match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
(Some(e), Some((hash, field))) if e.hash == hash => e
.min_effect
.map(|m| m.resolve(field, base.total.unwrap_or(metric.n)))
.unwrap_or(0.0),
_ => 0.0,
}
}
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,
})
}
pub(crate) struct BaselineRead {
pub value: f64,
pub kind: &'static str,
pub run_id: Option<String>,
pub best_before: Option<f64>,
pub total: Option<u64>,
pub cost: Option<f64>,
}
fn baseline_kind_for(
policy: &crate::policy::Policy,
metric: &crate::recommendation::MetricSnapshot,
) -> crate::policy::BaselineKind {
match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
(Some(e), Some((hash, _))) if e.hash == hash => e.baseline,
_ => crate::policy::BaselineKind::default(),
}
}
pub(crate) fn baseline_at_apply<S: SubstrateRead>(
sub: &S,
metric: &crate::recommendation::MetricSnapshot,
applied_at_ms: i64,
kind: crate::policy::BaselineKind,
cost_field: Option<&str>,
) -> Result<BaselineRead> {
use crate::policy::BaselineKind;
if let Some((evalset, field)) = crate::eval::parse_evalset_metric(&metric.metric) {
let before: Vec<(crate::eval::EvalRun, f64)> = crate::eval::eval_runs(sub, evalset, None)?
.into_iter()
.filter(|r| r.recorded_ms < applied_at_ms)
.filter_map(|r| crate::eval::run_value(&r, field).map(|v| (r, v)))
.collect();
if let Some(newest) = before.last() {
let best = before
.iter()
.fold(None::<&(crate::eval::EvalRun, f64)>, |acc, r| match acc {
None => Some(r),
Some(b) => {
let better = if metric.higher_is_better { r.1 > b.1 } else { r.1 < b.1 };
Some(if better { r } else { b })
}
})
.expect("non-empty");
let pick = match kind {
BaselineKind::NewestBeforeApply => newest,
BaselineKind::HighWater => best,
};
return Ok(BaselineRead {
value: pick.1,
kind: kind.as_str(),
run_id: Some(pick.0.run_id.clone()),
best_before: Some(best.1),
total: Some(pick.0.total()),
cost: cost_field.and_then(|f| crate::eval::run_value(&pick.0, f)),
});
}
}
Ok(BaselineRead { value: metric.baseline, kind: "snapshot", run_id: None, best_before: None, total: None, cost: None })
}
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 DecidedDraft {
grounded: bool,
sound: f64,
judged_by: crate::decide::JudgedBy,
}
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>>,
extra_statements: Vec<String>,
replay: Option<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 (depth, seg) in parents.iter().enumerate() {
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) => {
if depth == 0 && *seg == "retries" && !o.contains_key("retries") {
o.insert("retries".into(), Value::Object(serde_json::Map::new()));
}
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,
extra_statements: Vec::new(),
replay: 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,
extra_statements: Vec::new(),
replay: None,
})
}
P::Consolidation { lesson, supersedes } => {
let lesson = sanitize_lesson(&lesson);
if lesson.is_empty() {
return None;
}
let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("").to_string();
let mut hashes: Vec<String> = supersedes.into_iter().collect();
hashes.sort();
hashes.dedup();
let mut members = Vec::new();
for h in &hashes {
let g = sub.grain(h).ok().flatten()?;
if !g.is_live()
|| g.fact_relation() != Some("lesson")
|| g.fact_subject().is_none_or(|s| normalize_ident(s) != normalize_ident(&subject))
{
return None;
}
members.push(g);
}
if members.len() < 2 || members.iter().any(|m| m.namespace != members[0].namespace) {
return None;
}
let ns = members[0].namespace.clone();
let mut fields = fields;
if !ns.is_empty() {
fields.insert("namespace".into(), Value::from(ns.clone()));
}
fields.insert("consolidates".into(), Value::from(hashes.clone()));
let extra_statements: Vec<String> = members
.iter()
.map(|m| {
let mut marker = serde_json::Map::new();
marker.insert("subject".into(), Value::from(subject.clone()));
marker.insert("relation".into(), Value::from("mg:lesson_consolidated"));
marker.insert("object".into(), Value::from(lesson.clone()));
if !ns.is_empty() {
marker.insert("namespace".into(), Value::from(ns.clone()));
}
cal::supersede(&m.hash, "fact", &marker)
})
.collect();
args.insert("lesson".into(), Value::from(lesson.clone()));
args.insert("count".into(), Value::from(members.len() as u64));
Some(ResolvedProposal {
action: ActionKind::Consolidate,
proposal: Proposal::Cal { cal: cal::batch(&extra_statements) },
rendered: format!(
"Proposed consolidation of {} lessons on \"{subject}\" into one: \"{lesson}\"",
members.len()
),
summary_key: "llm.consolidation",
summary_args: args,
rollbackable: true,
evalset_hash: None,
importance: 0.6,
fact_fields: Some(fields),
extra_statements,
replay: 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),
extra_statements: Vec::new(),
replay: None,
})
}
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),
extra_statements: Vec::new(),
replay: None,
})
}
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,
extra_statements: Vec::new(),
replay: 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 replay = sub.plan_replay(hash, &body).ok().flatten();
let refused = match (&policy.plan_replay, &replay) {
(Some(gate), Some(report)) => gate.refusal(report),
_ => None,
};
let Value::Object(fields) = body else {
return None;
};
args.insert("plan".into(), Value::from(hash));
args.insert("edits".into(), Value::from(deltas.join("; ")));
if let Some(reason) = refused {
let mut data = serde_json::Map::new();
data.insert("plan".into(), Value::from(hash));
data.insert("edits".into(), Value::from(deltas.clone()));
data.insert("refused".into(), Value::from(reason.clone()));
args.insert("reason".into(), Value::from(reason.clone()));
return Some(ResolvedProposal {
action: ActionKind::Flag,
proposal: Proposal::Data { data },
rendered: format!(
"Plan revision ({}) refused by the rehearsal: {reason}",
deltas.join("; ")
),
summary_key: "llm.plan_revision_refused",
summary_args: args,
rollbackable: false,
evalset_hash: None,
importance: 0.4,
fact_fields: None,
extra_statements: Vec::new(),
replay,
});
}
let stmt = cal::supersede(hash, "workflow", &fields);
sub.validate_cal(&stmt).ok()?;
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,
extra_statements: Vec::new(),
replay,
})
}
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,
extra_statements: Vec::new(),
replay: None,
})
}
}
}
#[allow(clippy::too_many_arguments)]
fn stamp_llm(
model: &str,
d: &crate::llm::LlmDraft,
target_ref: String,
cited: Vec<String>,
resolved: Option<ResolvedProposal>,
confidence: f64,
now_ms: i64,
scope: &[String],
) -> 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 replay = resolved.as_ref().and_then(|r| r.replay.clone());
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)));
let mut statements = vec![cal::add("fact", &fields)];
statements.extend(r.extra_statements.iter().cloned());
r.proposal = Proposal::Cal { cal: cal::batch(&statements) };
}
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,
near_duplicate_of: Vec::new(),
replay,
scope: crate::recommendation::normalize_scope(scope),
judged_by: None,
llm_confidence: None,
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(),
baseline_kind: "snapshot".into(),
baseline_run_id: None,
best_before: None,
tolerance: 0.0,
current_run_id: None,
cost: None,
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,
baseline_kind: "snapshot".into(),
baseline_run_id: None,
best_before: None,
tolerance: 0.0,
current_run_id: None,
cost: None,
});
}
Ok(out)
}
fn moved_premises<S: OmsSubstrate>(
sub: &S,
rec: &Recommendation,
own: &[String],
) -> Result<(u64, u64)> {
let total = rec.evidence.len() as u64;
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;
}
}
}
}
}
}
Ok((moved, total))
}
fn withdraw_drifted_open<S: OmsSubstrate>(
sub: &mut S,
p: &mut LoopPersisted,
require_all: bool,
now_ms: i64,
) -> Result<u64> {
let open: Vec<String> = p
.status_index
.iter()
.filter(|(_, st)| matches!(st, RecStatus::Pending | RecStatus::Approved))
.map(|(h, _)| h.clone())
.collect();
let mut withdrawn = 0u64;
for rec_hash in open {
let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
if rec.evidence.is_empty() {
continue;
}
let (moved, total) = moved_premises(sub, &rec, &[])?;
let enough = if require_all { moved >= total } else { moved > 0 };
if moved == 0 || !enough {
continue;
}
let from = p.status_index.get(&rec_hash).copied().unwrap_or(RecStatus::Pending);
let prev = p.audit_heads.get(&rec_hash).cloned();
let audit = AuditRecord {
rec_hash: rec_hash.clone(),
from: Some(from),
to: RecStatus::Withdrawn,
actor: "engine:loop.premise_drift".into(),
observer_type: ObserverType::System,
because: format!(
"{moved} of {total} cited grains were superseded by a different value or retracted"
),
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, RecStatus::Withdrawn);
withdrawn += 1;
}
Ok(withdrawn)
}
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 latest_verdicts(p: &LoopPersisted) -> BTreeMap<String, String> {
let mut out = BTreeMap::new();
for (rec_hash, applied) in &p.applied {
let latest = p
.outcomes
.get(rec_hash)
.and_then(|v| v.iter().max_by_key(|o| o.measured_at_ms))
.map(|o| o.verdict.clone());
for h in &applied.created_hashes {
out.insert(h.clone(), latest.clone().unwrap_or_else(|| "unmeasured".into()));
}
}
out
}
pub const NEAR_DUPLICATE_COSINE: f64 = 0.90;
pub const NEAR_DUPLICATE_JACCARD: f64 = 0.60;
const NEAR_DUPLICATE_CAP: usize = 8;
pub(crate) fn near_duplicates_of<S: SubstrateRead + ?Sized>(
sub: &S,
subject: &str,
namespace: Option<&str>,
text: &str,
) -> Vec<crate::recommendation::NearDuplicate> {
use crate::analyzers::duplicate_sweep::{jaccard, tokenize};
if subject.is_empty() || text.trim().is_empty() {
return Vec::new();
}
let Ok(facts) = sub.grains_of_type(
crate::model::grain_type::FACT,
namespace,
ReadOpts { live_only: true, since_ms: None },
) else {
return Vec::new();
};
let mine = sub.embed(text).ok().flatten();
let my_tokens = tokenize(text);
let mut out: Vec<crate::recommendation::NearDuplicate> = facts
.iter()
.filter(|f| f.fact_relation() == Some("lesson"))
.filter(|f| f.fact_subject().is_some_and(|s| normalize_ident(s) == normalize_ident(subject)))
.filter_map(|f| {
let other = f.fact_object()?;
let (score, method, floor) = match (&mine, sub.embed(other).ok().flatten()) {
(Some(a), Some(b)) => (cosine(a, &b), "cosine", NEAR_DUPLICATE_COSINE),
_ => (jaccard(&my_tokens, &tokenize(other)), "jaccard", NEAR_DUPLICATE_JACCARD),
};
(score >= floor).then(|| crate::recommendation::NearDuplicate {
hash: f.hash.clone(),
score: (score * 1000.0).round() / 1000.0,
method: method.into(),
})
})
.collect();
out.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal).then(a.hash.cmp(&b.hash)));
out.truncate(NEAR_DUPLICATE_CAP);
out
}
fn cosine(a: &[f32], b: &[f32]) -> f64 {
if a.len() != b.len() || a.is_empty() {
return 0.0;
}
let (mut dot, mut na, mut nb) = (0f64, 0f64, 0f64);
for (x, y) in a.iter().zip(b) {
dot += *x as f64 * *y as f64;
na += *x as f64 * *x as f64;
nb += *y as f64 * *y as f64;
}
if na == 0.0 || nb == 0.0 {
0.0
} else {
dot / (na.sqrt() * nb.sqrt())
}
}
const CONSOLIDATION_INSTRUCTIONS: &str = " (8) {\"kind\":\"consolidation\",\"lesson\":\"...\",\
\"supersedes\":[\"<hash>\",...]} with the same entity target — ONLY in answer to a \
'Lesson pile' finding, which lists the live lessons on one entity that exceed \
its budget. Write ONE short imperative rule (max 240 chars) that says what \
those lessons say together, dropping nothing a lesson that measured 'held' \
required and keeping nothing only a lesson that measured 'regressed' or \
'drifted' added; 'supersedes' MUST be exactly the hashes that finding lists \
(cite them as evidence too). Applying it replaces every listed lesson with \
the one line; the reviewer can restore them all.";
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,
scope: &[String],
) -> 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,
near_duplicate_of: Vec::new(),
replay: None,
scope: crate::recommendation::normalize_scope(scope),
judged_by: d.judged_by,
llm_confidence: None,
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)
}
pub(crate) 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");
let mut q = plan();
q.as_object_mut().unwrap().remove("retries");
assert!(plan_set(&mut q, "retries.greet", json!(1)));
assert_eq!(plan_get(&q, "retries.greet"), json!(1));
let mut r = plan();
r.as_object_mut().unwrap().remove("edges");
assert!(!plan_set(&mut r, "edges.0.cond", json!("x")));
}
}
#[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""#));
}
}