1use crate::analyzer::{AnalyzeCtx, Analyzer, OutcomeInput};
11use crate::cal;
12use crate::config::{AppliedRecord, LoopPersisted};
13use crate::error::{Error, Result};
14use crate::manifest::{AnalyzerManifest, Capability};
15use crate::model::{normalize_ident, ActionKind, GrainRecord, Origin, Severity, TargetRef};
16use crate::recommendation::{
17 dedup_key, AuditRecord, ObserverType, Proposal, RecStatus, Recommendation, Summary,
18 MAX_BECAUSE, MAX_EVIDENCE,
19};
20use crate::substrate::{Capabilities, OmsSubstrate, ReadOpts, SubstrateRead};
21use serde::{Deserialize, Serialize};
22use serde_json::{Map, Value};
23use std::collections::{BTreeMap, BTreeSet};
24
25pub const LOOP_NS: &str = "areev-loop";
27
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
30pub enum Scope {
31 Read,
32 Write,
33 Review,
34 Apply,
35 Admin,
36}
37
38#[derive(Debug, Clone, Default)]
40pub struct ScopeSet(Vec<Scope>);
41
42impl ScopeSet {
43 pub fn of(scopes: &[Scope]) -> Self {
44 ScopeSet(scopes.to_vec())
45 }
46 pub fn all() -> Self {
49 ScopeSet(vec![Scope::Admin])
50 }
51 pub fn has(&self, s: Scope) -> bool {
52 self.0.contains(&Scope::Admin) || self.0.contains(&s)
53 }
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
58pub enum Decision {
59 Approve,
60 Reject,
61}
62
63#[derive(Debug, Clone, Default)]
65pub struct RunOptions {
66 pub min_new: Option<u64>,
67 pub min_new_errors: Option<u64>,
68 pub if_stale_ms: Option<i64>,
69 pub namespaces: Vec<String>,
71 pub full_sweep: bool,
78 pub triggering_actor: Option<String>,
85}
86
87#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(rename_all = "lowercase")]
90pub enum RunOutcome {
91 Ran,
92 Skipped,
93}
94
95#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
98#[serde(rename_all = "snake_case")]
99pub enum SkipReason {
100 MinNewNotMet,
101 NotStale,
102 LockHeld,
103}
104
105#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
107pub struct AnalyzerSkip {
108 pub id: String,
109 pub reason: String,
110}
111
112#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
114pub struct RunResult {
115 pub outcome: RunOutcome,
116 #[serde(skip_serializing_if = "Option::is_none")]
117 pub skip_reason: Option<SkipReason>,
118 pub new_grains: u64,
119 pub new_error_events: u64,
120 pub proposed: u64,
121 pub deduped: u64,
122 pub stored: u64,
123 #[serde(default)]
125 pub auto_applied: u64,
126 #[serde(default)]
127 pub analyzers_run: Vec<String>,
128 #[serde(default)]
129 pub analyzers_skipped: Vec<AnalyzerSkip>,
130}
131
132impl RunResult {
133 fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
134 RunResult {
135 outcome: RunOutcome::Skipped,
136 skip_reason: Some(reason),
137 new_grains,
138 new_error_events,
139 proposed: 0,
140 deduped: 0,
141 stored: 0,
142 auto_applied: 0,
143 analyzers_run: vec![],
144 analyzers_skipped: vec![],
145 }
146 }
147
148 pub fn ran(&self) -> bool {
149 self.outcome == RunOutcome::Ran
150 }
151}
152
153pub struct Engine {
156 analyzers: Vec<Box<dyn Analyzer>>,
157 policy: crate::policy::Policy,
158 llm: Option<Box<dyn crate::llm::LlmBackend>>,
161 ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
166}
167
168struct AnalysisPass {
169 survivors: Vec<Recommendation>,
170 proposed: u64,
171 deduped: u64,
172 analyzers_run: Vec<String>,
173 analyzers_skipped: Vec<AnalyzerSkip>,
174}
175
176impl Engine {
177 pub fn with_builtins() -> Self {
180 Engine {
181 analyzers: crate::analyzer::builtin_analyzers(),
182 policy: crate::policy::Policy::default(),
183 llm: None,
184 ground_llm: None,
185 }
186 }
187
188 pub fn empty() -> Self {
190 Engine {
191 analyzers: vec![],
192 policy: crate::policy::Policy::default(),
193 llm: None,
194 ground_llm: None,
195 }
196 }
197
198 pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
200 self.policy = policy;
201 self
202 }
203
204 pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
209 self.llm = Some(backend);
210 self
211 }
212
213 pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
217 self.ground_llm = Some(backend);
218 self
219 }
220
221 pub fn policy(&self) -> &crate::policy::Policy {
222 &self.policy
223 }
224
225 pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
227 self.analyzers.push(analyzer);
228 }
229
230 pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
231 &self.analyzers
232 }
233
234 pub fn analyze_only<S: OmsSubstrate>(
243 &self,
244 sub: &S,
245 opts: &RunOptions,
246 overrides: &BTreeMap<String, Map<String, Value>>,
247 now_ms: i64,
248 ) -> Result<Vec<Recommendation>> {
249 let persisted = LoopPersisted::from_value(sub.load_state()?)?;
250 let analysis_watermark = if opts.full_sweep {
251 None
252 } else {
253 persisted.state.watermark_ms
254 };
255 Ok(self
256 .analysis_pass(
257 sub,
258 &persisted,
259 opts,
260 overrides,
261 analysis_watermark,
262 now_ms,
263 &[],
264 )?
265 .survivors)
266 }
267
268 pub fn run<S: OmsSubstrate>(
271 &self,
272 sub: &mut S,
273 opts: &RunOptions,
274 now_ms: i64,
275 ) -> Result<RunResult> {
276 let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
277 let watermark = persisted.state.watermark_ms;
278 let analysis_watermark = if opts.full_sweep { None } else { watermark };
283
284 let (new_grains, new_error_events) = count_new(sub, watermark)?;
285 if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
286 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
287 }
288
289 let outcome_inputs = measure_outcomes(sub, &mut persisted, now_ms)?;
292
293 let AnalysisPass {
294 survivors,
295 proposed,
296 deduped,
297 analyzers_run,
298 analyzers_skipped,
299 } = self.analysis_pass(
300 &*sub,
301 &persisted,
302 opts,
303 &BTreeMap::new(),
304 analysis_watermark,
305 now_ms,
306 &outcome_inputs,
307 )?;
308
309 let mut stored = 0u64;
312 let mut auto_applied = 0u64;
313 for mut rec in survivors {
314 let spec = rec.to_grain_spec(LOOP_NS)?;
315 let hash = sub.put_grain(&spec)?;
316 rec.hash = hash.clone();
317 let actor = format!("engine:{}", rec.analyzer);
318 let audit = AuditRecord {
319 rec_hash: hash.clone(),
320 from: None,
321 to: RecStatus::Pending,
322 actor: actor.clone(),
323 observer_type: ObserverType::System,
324 because: "analyzer proposed".into(),
325 previous_audit_hash: None,
326 gating: None,
327 at_ms: now_ms,
328 };
329 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
330 persisted
331 .status_index
332 .insert(hash.clone(), RecStatus::Pending);
333 persisted.creators.insert(hash.clone(), actor);
334 if !matches!(rec.origin, Origin::Builtin) {
339 if let Some(trigger) = &opts.triggering_actor {
340 persisted.co_creators.insert(hash.clone(), trigger.clone());
341 }
342 }
343 persisted.audit_heads.insert(hash.clone(), audit_hash);
344 stored += 1;
345
346 if self.can_auto_apply(&*sub, &rec) {
347 self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
348 auto_applied += 1;
349 }
350 }
351
352 persisted.state.last_run_ms = Some(now_ms);
353 persisted.state.watermark_ms = Some(now_ms);
354 sub.store_state(&persisted.to_value()?)?;
355
356 Ok(RunResult {
357 outcome: RunOutcome::Ran,
358 skip_reason: None,
359 new_grains,
360 new_error_events,
361 proposed,
362 deduped,
363 stored,
364 auto_applied,
365 analyzers_run,
366 analyzers_skipped,
367 })
368 }
369
370 #[allow(clippy::too_many_arguments)]
374 fn analysis_pass<S: OmsSubstrate>(
375 &self,
376 sub: &S,
377 persisted: &LoopPersisted,
378 opts: &RunOptions,
379 external_overrides: &BTreeMap<String, Map<String, Value>>,
380 analysis_watermark: Option<i64>,
381 now_ms: i64,
382 outcome_inputs: &[OutcomeInput],
383 ) -> Result<AnalysisPass> {
384 let existing = existing_dedup_keys(sub, persisted)?;
385 let mut analyzers_run = Vec::new();
386 let mut analyzers_skipped = Vec::new();
387 let mut candidates: Vec<Recommendation> = Vec::new();
388 let caps = sub.capabilities();
389
390 for analyzer in &self.analyzers {
391 let m = analyzer.manifest();
392 let cfg = persisted.config.get(&m.id);
393 let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
394 if !enabled {
395 analyzers_skipped.push(AnalyzerSkip {
396 id: m.id.clone(),
397 reason: "disabled".into(),
398 });
399 continue;
400 }
401 if self.policy.denies(m.family()) {
402 analyzers_skipped.push(AnalyzerSkip {
403 id: m.id.clone(),
404 reason: "denied by host policy".into(),
405 });
406 continue;
407 }
408 if let Some(missing) = missing_capability(m, caps) {
409 analyzers_skipped.push(AnalyzerSkip {
410 id: m.id.clone(),
411 reason: format!("missing capability: {missing}"),
412 });
413 continue;
414 }
415 let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
416 if let Some(extra) = external_overrides.get(&m.id) {
417 for (key, value) in extra {
418 param_overrides.insert(key.clone(), value.clone());
419 }
420 }
421 let params = match m.resolve_params(¶m_overrides) {
422 Ok(p) => p,
423 Err(e) => {
424 analyzers_skipped.push(AnalyzerSkip {
425 id: m.id.clone(),
426 reason: e.to_string(),
427 });
428 continue;
429 }
430 };
431 let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
432 let ns_slice: &[String] = if ns_owned.is_empty() {
433 &opts.namespaces
434 } else {
435 &ns_owned
436 };
437 let reader: &dyn SubstrateRead = sub;
438 let ctx = AnalyzeCtx::new(
439 reader,
440 ¶ms,
441 ns_slice,
442 analysis_watermark,
443 now_ms,
444 outcome_inputs,
445 );
446 match analyzer.analyze(&ctx) {
447 Ok(drafts) => {
448 analyzers_run.push(m.id.clone());
449 for draft in drafts {
450 match stamp(m, ¶ms, draft, now_ms) {
451 Ok(rec) => candidates.push(rec),
452 Err(e) => analyzers_skipped.push(AnalyzerSkip {
453 id: m.id.clone(),
454 reason: e.to_string(),
455 }),
456 }
457 }
458 }
459 Err(e) => analyzers_skipped.push(AnalyzerSkip {
460 id: m.id.clone(),
461 reason: e.to_string(),
462 }),
463 }
464 }
465
466 if self.llm.is_some() {
467 candidates.extend(self.discover(
468 sub,
469 &candidates,
470 analysis_watermark,
471 &opts.namespaces,
472 now_ms,
473 ));
474 }
475
476 let proposed = candidates.len() as u64;
477 let mut seen = BTreeSet::new();
478 let mut survivors = Vec::new();
479 for candidate in candidates {
480 let family = crate::manifest::analyzer_family(&candidate.analyzer);
481 let floor = [
482 severity_floor_for(persisted, &candidate.analyzer),
483 self.policy.severity_floor(family),
484 ]
485 .into_iter()
486 .flatten()
487 .max();
488 if floor.is_some_and(|floor| candidate.severity < floor) {
489 continue;
490 }
491 if !seen.insert(candidate.dedup_key.clone()) {
492 continue;
493 }
494 if existing.contains(&candidate.dedup_key) {
495 continue;
496 }
497 if persisted
498 .cooldowns
499 .get(&candidate.dedup_key)
500 .is_some_and(|until| now_ms < *until)
501 {
502 continue;
503 }
504 survivors.push(candidate);
505 }
506 let deduped = proposed - survivors.len() as u64;
507 if self.llm.is_some() {
508 self.enrich(&mut survivors);
509 }
510 Ok(AnalysisPass {
511 survivors,
512 proposed,
513 deduped,
514 analyzers_run,
515 analyzers_skipped,
516 })
517 }
518
519 fn discover<S: OmsSubstrate>(
526 &self,
527 sub: &S,
528 candidates: &[Recommendation],
529 watermark: Option<i64>,
530 namespaces: &[String],
531 now_ms: i64,
532 ) -> Vec<Recommendation> {
533 let Some(llm) = &self.llm else {
534 return Vec::new();
535 };
536 let findings: Vec<crate::llm::FindingBrief> = candidates
537 .iter()
538 .take(32)
539 .map(|c| crate::llm::FindingBrief {
540 analyzer: c.analyzer.clone(),
541 summary: c.summary.render(),
542 target: c.target_ref.clone(),
543 severity: c.severity.as_str().to_string(),
544 })
545 .collect();
546 let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
552 let mut bundle: BTreeSet<String> = BTreeSet::new();
553 for c in candidates {
554 for h in &c.evidence {
555 if !bundle.contains(h) {
556 if let Ok(Some(g)) = sub.grain(h) {
557 push_evidence(&mut evidence, &mut bundle, &g);
558 }
559 }
560 }
561 }
562 let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
563 vec![None]
564 } else {
565 namespaces.iter().map(|n| Some(n.as_str())).collect()
566 };
567 let opts = ReadOpts { live_only: true, since_ms: watermark };
568 'seed: for gt in [
569 crate::model::grain_type::FACT,
570 crate::model::grain_type::OBSERVATION,
571 ] {
572 for ns in &scan_ns {
573 if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
574 for g in recent {
575 if evidence.len() >= 64 {
576 break 'seed;
577 }
578 push_evidence(&mut evidence, &mut bundle, &g);
579 }
580 }
581 }
582 }
583 if evidence.is_empty() {
584 return Vec::new(); }
586 let (approved, rejected) = self.llm_history(sub);
591 let request = crate::llm::LlmRequest {
592 loop_proto: 1,
593 op: "discover",
594 instructions: DISCOVER_INSTRUCTIONS,
595 findings: findings.clone(),
596 evidence: evidence.clone(),
597 rejected,
598 approved,
599 };
600 let Ok(body) = serde_json::to_string(&request) else {
601 return Vec::new();
602 };
603 let raw = match llm.complete(&body) {
604 Ok(r) => r,
605 Err(_) => return Vec::new(), };
607 let mut validated: Vec<(crate::llm::LlmDraft, String, Vec<String>)> = Vec::new();
611 for d in crate::llm::parse_discover(&raw)
612 .recommendations
613 .into_iter()
614 .take(crate::llm::MAX_LLM_DRAFTS)
615 {
616 let cited: Vec<String> =
617 d.evidence.iter().filter(|h| bundle.contains(*h)).cloned().collect();
618 if cited.is_empty() {
619 continue; }
621 let Ok(target) = TargetRef::parse(&d.target) else {
622 continue;
623 };
624 let tc = target.target_class();
625 if tc != "memory" && tc != "query" {
626 continue; }
628 validated.push((d, target.as_string(), cited));
629 }
630 if validated.is_empty() {
631 return Vec::new();
632 }
633 let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
640 self.verify_drafts(&**llm, ground, &validated, &evidence, now_ms)
641 }
642
643 fn verify_drafts(
650 &self,
651 llm: &dyn crate::llm::LlmBackend,
652 ground: &dyn crate::llm::LlmBackend,
653 validated: &[(crate::llm::LlmDraft, String, Vec<String>)],
654 evidence: &[crate::llm::EvidenceItem],
655 now_ms: i64,
656 ) -> Vec<Recommendation> {
657 use crate::llm::*;
658 let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
659 evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
660 let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
661 cited
662 .iter()
663 .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
664 .collect()
665 };
666
667 let claims: Vec<GroundItem> = validated
669 .iter()
670 .enumerate()
671 .map(|(i, (d, _t, cited))| GroundItem {
672 id: i,
673 claim: cap(&d.summary, MAX_SUMMARY_LEN),
674 evidence: ev_for(cited),
675 })
676 .collect();
677 let ground_req = GroundRequest {
678 loop_proto: 1,
679 op: "ground",
680 instructions: GROUND_INSTRUCTIONS,
681 claims,
682 };
683 let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
684 .ok()
685 .and_then(|b| ground.complete(&b).ok())
686 {
687 Some(raw) => parse_ground(&raw)
688 .results
689 .into_iter()
690 .filter(|r| r.supported)
691 .map(|r| r.id)
692 .collect(),
693 None => return Vec::new(),
694 };
695 if grounded.is_empty() {
696 return Vec::new();
697 }
698
699 let items: Vec<VerifyItem> = validated
705 .iter()
706 .enumerate()
707 .filter(|(i, _)| grounded.contains(i))
708 .map(|(i, (d, t, cited))| VerifyItem {
709 id: i,
710 summary: cap(&d.summary, MAX_SUMMARY_LEN),
711 target: t.clone(),
712 evidence: ev_for(cited),
713 })
714 .collect();
715 let verify_req = VerifyRequest {
716 loop_proto: 1,
717 op: "verify",
718 instructions: VERIFY_INSTRUCTIONS,
719 findings: items,
720 };
721 let verdicts: std::collections::BTreeMap<usize, f64> =
722 match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
723 Some(raw) => parse_verify(&raw)
724 .results
725 .into_iter()
726 .filter(|r| r.keep)
727 .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
728 .collect(),
729 None => return Vec::new(),
730 };
731
732 let mut out = Vec::new();
736 for (i, (d, target_str, cited)) in validated.iter().enumerate() {
737 if let Some(&conf) = verdicts.get(&i) {
738 if conf >= MIN_LLM_CONFIDENCE {
739 out.push(stamp_llm(
740 llm.model(),
741 d,
742 target_str.clone(),
743 cited.clone(),
744 conf,
745 now_ms,
746 ));
747 }
748 }
749 }
750 out
751 }
752
753 fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
758 const MAX: usize = 20;
759 let Ok(mut recs) = self.recommendations(sub, None) else {
760 return (Vec::new(), Vec::new());
761 };
762 recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
763 recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
764 let mut approved = Vec::new();
765 let mut rejected = Vec::new();
766 for r in &recs {
767 match r.status {
768 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
769 if approved.len() < MAX =>
770 {
771 approved.push(r.summary.render());
772 }
773 RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
774 _ => {}
775 }
776 }
777 (approved, rejected)
778 }
779
780 fn enrich(&self, survivors: &mut [Recommendation]) {
785 let Some(llm) = &self.llm else {
786 return;
787 };
788 if survivors.is_empty() {
789 return;
790 }
791 let findings: Vec<crate::llm::FindingBrief> = survivors
792 .iter()
793 .map(|r| crate::llm::FindingBrief {
794 analyzer: r.analyzer.clone(),
795 summary: r.summary.render(),
796 target: r.target_ref.clone(),
797 severity: r.severity.as_str().to_string(),
798 })
799 .collect();
800 let request = crate::llm::LlmRequest {
801 loop_proto: 1,
802 op: "enrich",
803 instructions: ENRICH_INSTRUCTIONS,
804 findings,
805 evidence: Vec::new(),
806 rejected: Vec::new(),
807 approved: Vec::new(),
808 };
809 let Ok(body) = serde_json::to_string(&request) else {
810 return;
811 };
812 let raw = match llm.complete(&body) {
813 Ok(r) => r,
814 Err(_) => return,
815 };
816 for note in crate::llm::parse_enrich(&raw).notes {
817 if note.guidance.trim().is_empty() {
818 continue;
819 }
820 if let Some(r) = survivors
821 .iter_mut()
822 .find(|r| r.target_ref == note.target && r.guidance.is_none())
823 {
824 r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
825 }
826 }
827 }
828
829 fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
838 if !rec.origin.auto_apply_eligible() || rec.destructive {
839 return false;
840 }
841 let manifest_ok = self
845 .analyzers
846 .iter()
847 .map(|a| a.manifest())
848 .find(|m| m.id == rec.analyzer)
849 .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
850 if !manifest_ok {
851 return false;
852 }
853 let Ok(target) = TargetRef::parse(&rec.target_ref) else {
854 return false;
855 };
856 let family = crate::manifest::analyzer_family(&rec.analyzer);
857 if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
858 return false;
859 }
860 match &rec.proposal {
865 Proposal::Cal { cal } => cal
866 .lines()
867 .map(str::trim)
868 .filter(|l| !l.is_empty())
869 .all(|l| supersede_is_value_identical(sub, l)),
870 _ => false,
871 }
872 }
873
874 fn auto_apply<S: OmsSubstrate>(
877 &self,
878 sub: &mut S,
879 p: &mut LoopPersisted,
880 rec: &Recommendation,
881 now_ms: i64,
882 ) -> Result<()> {
883 let mut created = Vec::new();
884 if let Proposal::Cal { cal } = &rec.proposal {
885 if cal.lines().map(str::trim).any(is_definition_statement) {
890 return Err(Error::InvalidProposal(
891 "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
892 auto-applied: it changes what every future context contains, so it \
893 requires a human APPROVE + APPLY with BECAUSE"
894 .into(),
895 ));
896 }
897 for r in sub.execute_cal(cal)? {
898 if let Some(h) = r.get("hash").and_then(Value::as_str) {
899 created.push(h.to_string());
900 }
901 }
902 }
903 let applied = AppliedRecord {
904 applied_at_ms: now_ms,
905 target_ref: rec.target_ref.clone(),
906 rollbackable: rec.rollbackable,
907 created_hashes: created,
908 inverse_cal: None,
909 metric: rec.metric.clone(),
910 };
911 let prev = p.audit_heads.get(&rec.hash).cloned();
912 let audit = AuditRecord {
913 rec_hash: rec.hash.clone(),
914 from: Some(RecStatus::Pending),
915 to: RecStatus::Applied,
916 actor: "policy:auto".into(),
917 observer_type: ObserverType::Policy,
918 because: "auto-applied per host policy".into(),
919 previous_audit_hash: prev,
920 gating: None,
921 at_ms: now_ms,
922 };
923 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
924 p.audit_heads.insert(rec.hash.clone(), audit_hash);
925 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
926 p.applied.insert(rec.hash.clone(), applied);
927 Ok(())
928 }
929
930 #[allow(clippy::too_many_arguments)]
933 pub fn review<S: OmsSubstrate>(
934 &self,
935 sub: &mut S,
936 rec_hash: &str,
937 decision: Decision,
938 actor: &str,
939 observer: ObserverType,
940 scopes: &ScopeSet,
941 because: &str,
942 now_ms: i64,
943 ) -> Result<()> {
944 if !scopes.has(Scope::Review) {
945 return Err(Error::ScopeDenied("review".into()));
946 }
947 let because = validate_because(because)?;
948 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
949 let status = *p
950 .status_index
951 .get(rec_hash)
952 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
953 let to = match decision {
954 Decision::Approve => RecStatus::Approved,
955 Decision::Reject => RecStatus::Rejected,
956 };
957 if !status.can_transition_to(to, false) {
958 return Err(Error::LifecycleViolation(format!(
959 "{} -> {}",
960 status.as_str(),
961 to.as_str()
962 )));
963 }
964 if to == RecStatus::Approved {
965 if let Some(creator) = p.creators.get(rec_hash) {
966 if creator == actor {
967 return Err(Error::SelfApproval(format!(
968 "{actor} created this recommendation"
969 )));
970 }
971 }
972 if let Some(trigger) = p.co_creators.get(rec_hash) {
973 if trigger == actor {
974 return Err(Error::SelfApproval(format!(
975 "{actor} triggered the run that authored this recommendation"
976 )));
977 }
978 }
979 }
980 let prev = p.audit_heads.get(rec_hash).cloned();
981 let audit = AuditRecord {
982 rec_hash: rec_hash.into(),
983 from: Some(status),
984 to,
985 actor: actor.into(),
986 observer_type: observer,
987 because,
988 previous_audit_hash: prev,
989 gating: None,
990 at_ms: now_ms,
991 };
992 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
993 p.audit_heads.insert(rec_hash.into(), audit_hash);
994 p.status_index.insert(rec_hash.into(), to);
995 if to == RecStatus::Rejected {
996 if let Ok(rec) = load_rec(sub, rec_hash) {
997 const BASE_MS: i64 = 7 * 86_400_000;
1002 const CAP_MS: i64 = 90 * 86_400_000;
1003 let strikes = p.cooldown_strikes.entry(rec.dedup_key.clone()).or_insert(0);
1004 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
1005 *strikes = strikes.saturating_add(1);
1006 p.cooldowns.insert(rec.dedup_key, now_ms + interval);
1007 }
1008 }
1009 sub.store_state(&p.to_value()?)?;
1010 Ok(())
1011 }
1012
1013 pub fn preflight_apply<S: OmsSubstrate>(
1030 &self,
1031 sub: &S,
1032 rec_hash: &str,
1033 scopes: &ScopeSet,
1034 allow_destructive: bool,
1035 has_gating: bool,
1036 ) -> Result<()> {
1037 if !scopes.has(Scope::Apply) {
1038 return Err(Error::ScopeDenied("apply".into()));
1039 }
1040 let rec = load_rec(sub, rec_hash)?;
1041 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1042 return Err(Error::DestructiveGated(
1043 "destructive apply requires admin scope + allow_destructive".into(),
1044 ));
1045 }
1046 ensure_executable(rec.action_kind, &rec.proposal)?;
1047 if requires_gating(rec.action_kind) && !has_gating {
1048 return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1049 }
1050 Ok(())
1051 }
1052
1053 #[allow(clippy::too_many_arguments)]
1057 pub fn apply<S: OmsSubstrate>(
1058 &self,
1059 sub: &mut S,
1060 rec_hash: &str,
1061 actor: &str,
1062 observer: ObserverType,
1063 scopes: &ScopeSet,
1064 because: &str,
1065 allow_destructive: bool,
1066 now_ms: i64,
1067 ) -> Result<AppliedRecord> {
1068 self.apply_inner(
1069 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1070 )
1071 }
1072
1073 pub fn gating_evidence<S: OmsSubstrate>(
1080 &self,
1081 sub: &S,
1082 rec_hash: &str,
1083 run_id: &str,
1084 ) -> Result<crate::recommendation::GatingEvidence> {
1085 let rec = self
1086 .recommendations(sub, None)?
1087 .into_iter()
1088 .find(|r| r.hash == rec_hash)
1089 .ok_or_else(|| {
1090 Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1091 })?;
1092 let pin = rec.evalset_hash.ok_or_else(|| {
1093 Error::InvalidProposal(
1094 "this recommendation pins no evalset — a gating run applies only \
1095 to code and adapter revisions"
1096 .into(),
1097 )
1098 })?;
1099 match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1104 Some(run) => Ok(crate::recommendation::GatingEvidence {
1105 evalset_hash: pin,
1106 run_id: run.run_id,
1107 passed: run.passed,
1108 failed: run.failed,
1109 }),
1110 None => Err(Error::InvalidProposal(format!(
1111 "no recorded gate run '{run_id}' for evalset {pin} — run \
1112 `areev eval run --evalset {pin} ...` first"
1113 ))),
1114 }
1115 }
1116
1117 #[allow(clippy::too_many_arguments)]
1121 pub fn apply_gated<S: OmsSubstrate>(
1122 &self,
1123 sub: &mut S,
1124 rec_hash: &str,
1125 actor: &str,
1126 observer: ObserverType,
1127 scopes: &ScopeSet,
1128 because: &str,
1129 allow_destructive: bool,
1130 gating: &crate::recommendation::GatingEvidence,
1131 now_ms: i64,
1132 ) -> Result<AppliedRecord> {
1133 self.apply_inner(
1134 sub,
1135 rec_hash,
1136 actor,
1137 observer,
1138 scopes,
1139 because,
1140 allow_destructive,
1141 Some(gating),
1142 now_ms,
1143 )
1144 }
1145
1146 #[allow(clippy::too_many_arguments)]
1147 fn apply_inner<S: OmsSubstrate>(
1148 &self,
1149 sub: &mut S,
1150 rec_hash: &str,
1151 actor: &str,
1152 observer: ObserverType,
1153 scopes: &ScopeSet,
1154 because: &str,
1155 allow_destructive: bool,
1156 gating: Option<&crate::recommendation::GatingEvidence>,
1157 now_ms: i64,
1158 ) -> Result<AppliedRecord> {
1159 if !scopes.has(Scope::Apply) {
1160 return Err(Error::ScopeDenied("apply".into()));
1161 }
1162 let because = validate_because(because)?;
1163 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1164 let status = *p
1165 .status_index
1166 .get(rec_hash)
1167 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1168 if !status.can_transition_to(RecStatus::Applied, false) {
1169 return Err(Error::LifecycleViolation(format!(
1170 "{} -> applied (approve first)",
1171 status.as_str()
1172 )));
1173 }
1174 let rec = load_rec(sub, rec_hash)?;
1175 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1176 return Err(Error::DestructiveGated(
1177 "destructive apply requires admin scope + allow_destructive".into(),
1178 ));
1179 }
1180 if requires_gating(rec.action_kind) {
1185 let g = gating
1186 .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1187 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1188 if g.evalset_hash != pin {
1189 return Err(Error::InvalidProposal(format!(
1190 "gating ran evalset {} but the recommendation is pinned \
1191 to {pin} (Rule E1)",
1192 g.evalset_hash
1193 )));
1194 }
1195 match sub.grain(pin)? {
1196 Some(evalset) if evalset.is_live() => {}
1197 Some(_) => {
1198 return Err(Error::InvalidProposal(
1199 "the pinned evalset was superseded after gating — \
1200 the recommendation must re-gate (Rule E1)"
1201 .into(),
1202 ))
1203 }
1204 None => {
1205 return Err(Error::InvalidProposal(format!(
1206 "pinned evalset {pin} not found in the substrate"
1207 )))
1208 }
1209 }
1210 if g.failed > 0 {
1211 return Err(Error::InvalidProposal(format!(
1212 "the gating run failed {}/{} cases — a failing gate \
1213 admits nothing",
1214 g.failed,
1215 g.passed + g.failed
1216 )));
1217 }
1218 }
1219
1220 let mut created = Vec::new();
1222 let mut inverse_cal: Option<String> = None;
1226 match &rec.proposal {
1227 Proposal::Cal { cal } => {
1228 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1229 if !is_definition_statement(line) {
1230 continue;
1231 }
1232 match sub.definition_inverse(line)? {
1233 Some(inv) => inverse_cal = Some(inv),
1234 None => {
1235 return Err(Error::InvalidProposal(format!(
1236 "this substrate cannot record a rollback inverse for {line:?}; \
1237 a definition rewrite that ROLLBACK could not undo is refused \
1238 rather than applied"
1239 )))
1240 }
1241 }
1242 }
1243 let rows = sub.execute_cal(cal)?;
1244 for r in rows {
1245 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1246 created.push(h.to_string());
1247 }
1248 }
1249 }
1250 Proposal::Edit { .. } => {
1253 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1254 }
1255 Proposal::Data { data } if requires_gating(rec.action_kind) => {
1263 let relation = if rec.action_kind == ActionKind::AdapterRevision {
1264 "mg:adapter_promotion"
1265 } else {
1266 "mg:code_promotion"
1267 };
1268 let mut spec = crate::substrate::GrainSpec::new(
1269 crate::model::grain_type::FACT,
1270 LOOP_NS,
1271 )
1272 .with_field("subject", rec.target_ref.clone())
1273 .with_field("relation", relation)
1274 .with_field(
1275 "object",
1276 serde_json::to_string(data)
1277 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1278 )
1279 .with_field("rec_hash", rec_hash.to_string());
1280 if let Some(g) = gating {
1281 spec = spec
1282 .with_field("gating_evalset", g.evalset_hash.clone())
1283 .with_field("gating_run_id", g.run_id.clone());
1284 }
1285 created.push(sub.put_grain(&spec)?);
1286 }
1287 Proposal::Data { data } => {
1288 let revert_of = data
1293 .get("revert_of")
1294 .and_then(Value::as_str)
1295 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1296 self.rollback(
1297 sub,
1298 revert_of,
1299 actor,
1300 observer,
1301 scopes,
1302 &because,
1303 now_ms,
1304 )?;
1305 p = LoopPersisted::from_value(sub.load_state()?)?;
1308 }
1309 }
1310
1311 let applied = AppliedRecord {
1312 applied_at_ms: now_ms,
1313 target_ref: rec.target_ref.clone(),
1314 rollbackable: rec.rollbackable,
1315 created_hashes: created,
1316 inverse_cal,
1317 metric: rec.metric.clone(),
1318 };
1319 let prev = p.audit_heads.get(rec_hash).cloned();
1320 let audit = AuditRecord {
1321 rec_hash: rec_hash.into(),
1322 from: Some(status),
1323 to: RecStatus::Applied,
1324 actor: actor.into(),
1325 observer_type: observer,
1326 because,
1327 previous_audit_hash: prev,
1328 gating: gating.cloned(),
1329 at_ms: now_ms,
1330 };
1331 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1332 p.audit_heads.insert(rec_hash.into(), audit_hash);
1333 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1334 p.applied.insert(rec_hash.into(), applied.clone());
1335 sub.store_state(&p.to_value()?)?;
1336 Ok(applied)
1337 }
1338
1339 #[allow(clippy::too_many_arguments)]
1342 pub fn rollback<S: OmsSubstrate>(
1343 &self,
1344 sub: &mut S,
1345 rec_hash: &str,
1346 actor: &str,
1347 observer: ObserverType,
1348 scopes: &ScopeSet,
1349 because: &str,
1350 now_ms: i64,
1351 ) -> Result<()> {
1352 if !scopes.has(Scope::Apply) {
1353 return Err(Error::ScopeDenied("apply".into()));
1354 }
1355 let because = validate_because(because)?;
1356 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1357 let status = *p
1358 .status_index
1359 .get(rec_hash)
1360 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1361 if !status.can_transition_to(RecStatus::RolledBack, false) {
1362 return Err(Error::LifecycleViolation(format!(
1363 "{} -> rolled_back",
1364 status.as_str()
1365 )));
1366 }
1367 let applied = p
1368 .applied
1369 .get(rec_hash)
1370 .cloned()
1371 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1372 if !applied.rollbackable {
1373 return Err(Error::LifecycleViolation(
1374 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1375 ));
1376 }
1377 for h in &applied.created_hashes {
1378 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1379 }
1380 if let Some(inverse) = &applied.inverse_cal {
1387 sub.execute_cal(inverse)?;
1388 }
1389 let prev = p.audit_heads.get(rec_hash).cloned();
1390 let audit = AuditRecord {
1391 rec_hash: rec_hash.into(),
1392 from: Some(status),
1393 to: RecStatus::RolledBack,
1394 actor: actor.into(),
1395 observer_type: observer,
1396 because,
1397 previous_audit_hash: prev,
1398 gating: None,
1399 at_ms: now_ms,
1400 };
1401 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1402 p.audit_heads.insert(rec_hash.into(), audit_hash);
1403 p.status_index
1404 .insert(rec_hash.into(), RecStatus::RolledBack);
1405 sub.store_state(&p.to_value()?)?;
1406 Ok(())
1407 }
1408
1409 pub fn recommendations<S: OmsSubstrate>(
1414 &self,
1415 sub: &S,
1416 status_filter: Option<RecStatus>,
1417 ) -> Result<Vec<Recommendation>> {
1418 let p = LoopPersisted::from_value(sub.load_state()?)?;
1419 let grains = sub.grains_of_type(
1420 crate::model::grain_type::RECOMMENDATION,
1421 Some(LOOP_NS),
1422 ReadOpts {
1423 live_only: false,
1424 since_ms: None,
1425 },
1426 )?;
1427 let mut out = Vec::new();
1428 for g in grains {
1429 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1430 rec.status = p
1431 .status_index
1432 .get(&g.hash)
1433 .copied()
1434 .unwrap_or(RecStatus::Pending);
1435 if let Some(f) = status_filter {
1436 if rec.status != f {
1437 continue;
1438 }
1439 }
1440 out.push(rec);
1441 }
1442 out.sort_by(|a, b| {
1451 b.severity
1452 .cmp(&a.severity)
1453 .then(a.created_at_ms.cmp(&b.created_at_ms))
1454 .then(a.dedup_key.cmp(&b.dedup_key))
1455 .then(a.hash.cmp(&b.hash))
1456 });
1457 Ok(out)
1458 }
1459
1460 pub fn analyzer_settings<S: OmsSubstrate>(
1463 &self,
1464 sub: &S,
1465 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1466 let p = LoopPersisted::from_value(sub.load_state()?)?;
1467 Ok(self
1468 .analyzers
1469 .iter()
1470 .map(|a| {
1471 let m = a.manifest();
1472 let cfg = p.config.get(&m.id);
1473 crate::config::AnalyzerSetting {
1474 id: m.id.clone(),
1475 title: m.title.clone(),
1476 description: m.description.clone(),
1477 tier: format!("{:?}", m.tier),
1478 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1479 default_on: m.default_on,
1480 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1481 severity_floor: cfg
1482 .and_then(|c| c.severity_floor)
1483 .map(|s| s.as_str().to_string()),
1484 }
1485 })
1486 .collect())
1487 }
1488
1489 pub fn set_analyzer_config<S: OmsSubstrate>(
1495 &self,
1496 sub: &mut S,
1497 analyzer_id: &str,
1498 update: crate::config::AnalyzerConfigUpdate,
1499 scopes: &ScopeSet,
1500 ) -> Result<crate::config::AnalyzerConfig> {
1501 if !scopes.has(Scope::Admin) {
1502 return Err(Error::ScopeDenied("admin".into()));
1503 }
1504 let manifest = self
1505 .analyzers
1506 .iter()
1507 .map(|a| a.manifest())
1508 .find(|m| m.id == analyzer_id)
1509 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1510 if let Some(params) = &update.params {
1512 manifest.resolve_params(params)?;
1513 }
1514 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1515 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1516 if let Some(enabled) = update.enabled {
1517 cfg.enabled = Some(enabled);
1518 }
1519 if update.clear_floor {
1520 cfg.severity_floor = None;
1521 } else if let Some(floor) = update.severity_floor {
1522 cfg.severity_floor = Some(floor);
1523 }
1524 if let Some(params) = update.params {
1525 cfg.params = params;
1526 }
1527 if let Some(ns) = update.namespaces {
1528 cfg.namespaces = ns;
1529 }
1530 let stored = cfg.clone();
1531 sub.store_state(&p.to_value()?)?;
1532 Ok(stored)
1533 }
1534
1535 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1538 let p = LoopPersisted::from_value(sub.load_state()?)?;
1539 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1540 out.sort_by(|a, b| {
1544 a.measured_at_ms
1545 .cmp(&b.measured_at_ms)
1546 .then(a.horizon_ms.cmp(&b.horizon_ms))
1547 .then(a.metric.cmp(&b.metric))
1548 .then(a.rec_hash.cmp(&b.rec_hash))
1549 });
1550 Ok(out)
1551 }
1552
1553 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1557 let p = LoopPersisted::from_value(sub.load_state()?)?;
1558 let (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
1559 let recs = self.recommendations(sub, None)?;
1560 let mut pending = 0;
1561 let mut applied = 0;
1562 for r in &recs {
1563 match r.status {
1564 RecStatus::Pending => pending += 1,
1565 RecStatus::Applied => applied += 1,
1566 _ => {}
1567 }
1568 }
1569 let stale = match p.state.last_run_ms {
1571 None => true,
1572 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1573 };
1574 Ok(Health {
1575 last_run_ms: p.state.last_run_ms,
1576 grains_since_run,
1577 error_events_since_run,
1578 pending,
1579 applied,
1580 total: recs.len() as u64,
1581 stale,
1582 })
1583 }
1584
1585 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1590 let recs = self.recommendations(sub, None)?;
1591 let mut m = LlmMetrics::default();
1592 for r in &recs {
1593 if !matches!(r.origin, Origin::Llm { .. }) {
1594 continue;
1595 }
1596 m.proposed += 1;
1597 match r.status {
1598 RecStatus::Pending => m.pending += 1,
1599 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1600 RecStatus::Rejected => m.rejected += 1,
1601 RecStatus::Expired => {}
1602 }
1603 }
1604 let decided = m.approved + m.rejected;
1605 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1606 Ok(m)
1607 }
1608}
1609
1610#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1612pub struct Health {
1613 #[serde(skip_serializing_if = "Option::is_none")]
1614 pub last_run_ms: Option<i64>,
1615 pub grains_since_run: u64,
1616 pub error_events_since_run: u64,
1617 pub pending: u64,
1618 pub applied: u64,
1619 pub total: u64,
1620 pub stale: bool,
1623}
1624
1625#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1627pub struct LlmMetrics {
1628 pub proposed: u64,
1631 pub pending: u64,
1632 pub approved: u64,
1634 pub rejected: u64,
1635 #[serde(skip_serializing_if = "Option::is_none")]
1637 pub approval_rate: Option<f64>,
1638}
1639
1640fn measure_outcomes<S: OmsSubstrate>(
1647 sub: &S,
1648 p: &mut LoopPersisted,
1649 now_ms: i64,
1650) -> Result<Vec<OutcomeInput>> {
1651 let mut due: Vec<(String, crate::config::AppliedRecord, i64)> = Vec::new();
1653 for (h, a) in &p.applied {
1654 if p.status_index.get(h) != Some(&RecStatus::Applied) {
1655 continue;
1656 }
1657 let Some(metric) = &a.metric else { continue };
1658 let done = p.measured.get(h).cloned().unwrap_or_default();
1659 for horizon in metric.horizons() {
1660 if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
1661 due.push((h.clone(), a.clone(), horizon));
1662 }
1663 }
1664 }
1665
1666 let mut out = Vec::new();
1667 for (rec_hash, applied, horizon) in due {
1668 let metric = applied.metric.as_ref().unwrap();
1669 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1670 continue; };
1672 let regressed = crate::recommendation::is_regression(
1673 metric.baseline,
1674 current,
1675 metric.higher_is_better,
1676 );
1677 p.outcomes.entry(rec_hash.clone()).or_default().push(
1678 crate::recommendation::OutcomeResult {
1679 rec_hash: rec_hash.clone(),
1680 metric: metric.metric.clone(),
1681 baseline: metric.baseline,
1682 current,
1683 verdict: if regressed { "regressed" } else { "held" }.into(),
1684 horizon_ms: horizon,
1685 measured_at_ms: now_ms,
1686 },
1687 );
1688 p.measured.entry(rec_hash.clone()).or_default().push(horizon);
1689 if regressed {
1690 out.push(OutcomeInput {
1691 rec_hash,
1692 target_ref: applied.target_ref.clone(),
1693 metric: metric.metric.clone(),
1694 baseline: metric.baseline,
1695 current,
1696 unit: metric.unit.clone(),
1697 higher_is_better: metric.higher_is_better,
1698 });
1699 }
1700 }
1701 Ok(out)
1702}
1703
1704pub(crate) fn measure_metric<S: SubstrateRead>(
1706 sub: &S,
1707 metric: &crate::recommendation::MetricSnapshot,
1708 since_ms: i64,
1709) -> Result<Option<f64>> {
1710 match metric.metric.as_str() {
1711 "tool_error_recurrence" => {
1716 let Some(tool) = &metric.subject else { return Ok(None) };
1717 let tools = sub.grains_of_type(
1718 crate::model::grain_type::TOOL,
1719 None,
1720 ReadOpts { live_only: true, since_ms: Some(since_ms) },
1721 )?;
1722 let n = tools
1723 .iter()
1724 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
1725 .filter(|t| {
1726 metric.relation.as_deref().is_none_or(|sig| {
1729 crate::analyzers::tool_failure::normalize_signature(
1730 t.tool_content().unwrap_or(""),
1731 ) == sig
1732 })
1733 })
1734 .count();
1735 Ok(Some(n as f64))
1736 }
1737 "contradiction_recurrence" => {
1741 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
1742 return Ok(None);
1743 };
1744 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
1745 let distinct: BTreeSet<String> = facts
1746 .iter()
1747 .filter(|f| {
1748 f.fact_relation()
1749 .is_some_and(|r| normalize_ident(r) == *relation)
1750 })
1751 .filter_map(|f| f.fact_object().map(normalize_ident))
1752 .collect();
1753 Ok(Some(distinct.len().saturating_sub(1) as f64))
1754 }
1755 m if m.starts_with("evalset:") => {
1768 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
1769 return Ok(None);
1770 };
1771 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
1772 return Ok(None);
1773 };
1774 Ok(match field {
1778 "failed" => Some(run.failed as f64),
1779 "passed" => Some(run.passed as f64),
1780 "total" => Some(run.total() as f64),
1781 "error_rate" => match run.total() {
1782 0 => None, t => Some(run.failed as f64 / t as f64),
1784 },
1785 other => run.field(other),
1786 })
1787 }
1788 _ => Ok(None),
1789 }
1790}
1791
1792fn scoped_live_facts<S: SubstrateRead>(
1795 sub: &S,
1796 namespace: Option<&str>,
1797 subject: &str,
1798) -> Result<Vec<GrainRecord>> {
1799 let facts = sub.grains_of_type(
1800 crate::model::grain_type::FACT,
1801 None,
1802 ReadOpts { live_only: true, since_ms: None },
1803 )?;
1804 Ok(facts
1805 .into_iter()
1806 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
1807 .filter(|f| {
1808 f.fact_subject()
1809 .is_some_and(|s| normalize_ident(s) == subject)
1810 })
1811 .collect())
1812}
1813
1814fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
1826 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
1827 return false;
1828 };
1829 if fields.is_empty() {
1830 return false;
1831 }
1832 let Ok(Some(grain)) = sub.grain(&target) else {
1833 return false;
1834 };
1835 if grain.valid_to_ms.is_some() {
1845 return false;
1846 }
1847 fields.iter().all(|(k, v)| {
1849 if k == "namespace" {
1850 return v
1851 .as_str()
1852 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
1853 }
1854 match (v, grain.fields.get(k)) {
1855 (Value::String(a), Some(Value::String(b))) => {
1856 normalize_ident(a) == normalize_ident(b)
1857 }
1858 (a, Some(b)) => a == b,
1859 (_, None) => false,
1860 }
1861 })
1862}
1863
1864const MIN_LLM_CONFIDENCE: f64 = 0.75;
1867
1868const DISCOVER_INSTRUCTIONS: &str = "You review an agent's memory for quality. \
1873Given deterministic findings and the evidence they cite, propose ADDITIONAL \
1874findings the deterministic checks would miss (e.g. a semantic contradiction, a \
1875stale assumption, a duplicated meaning). SCORING: propose a finding ONLY if you \
1876are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
1877useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
1878earns 0. When in doubt, propose nothing — an empty list is the correct answer \
1879when there is nothing worth flagging. The 'approved' and 'rejected' lists, when \
1880present, show findings this reviewer recently accepted or rejected — prefer the \
1881kind they accept and avoid the kind they reject. Every proposal MUST cite one \
1882or more evidence hashes from the bundle, target a memory entity, and include \
1883your confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
1884\"target\":\"entity:<ns>/<subject>\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
1885\"confidence\":0.0}]}. Propose nothing you cannot ground in the evidence.";
1886
1887const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
1893fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
1894the facts it relies on are actually present in the cited evidence, NOT that its \
1895conclusion is stated verbatim. Decompose the finding into the factual claims it \
1896depends on. Mark supported=true when those facts are present in the evidence \
1897(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
1898on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
1899different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
1900
1901const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
1903each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
1904never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
1905SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
1906'possible' findings with no concrete defect, and reject any claimed \
1907inconsistency or contradiction that is not backed by at least two actually \
1908conflicting facts in the cited evidence. (2) Context — does the finding \
1909correctly read its cited evidence, or misinterpret what the grains say? Keep a \
1910finding when it names a genuine, specific problem grounded in its evidence and \
1911materially useful to a human reviewer; otherwise reject it, and default to \
1912keep=false when uncertain. Do NOT reject a finding for being 'already known', \
1913redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
1914grounded cross-fact inconsistency with two conflicting facts is exactly what to \
1915KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
1916{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
1917
1918const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
1920guidance note to help a human reviewer decide. Do not restate the finding. Return \
1921JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
1922
1923fn push_evidence(
1926 evidence: &mut Vec<crate::llm::EvidenceItem>,
1927 bundle: &mut BTreeSet<String>,
1928 g: &GrainRecord,
1929) {
1930 if evidence.len() < 64 && bundle.insert(g.hash.clone()) {
1931 evidence.push(crate::llm::EvidenceItem {
1932 hash: g.hash.clone(),
1933 grain_type: g.grain_type.clone(),
1934 text: crate::llm::cap(&grain_brief(g), 400),
1935 });
1936 }
1937}
1938
1939fn grain_brief(g: &GrainRecord) -> String {
1941 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
1942 return format!("{s} {r} {o}");
1943 }
1944 for key in ["content", "body", "text", "summary"] {
1945 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
1946 if !v.is_empty() {
1947 return v.to_string();
1948 }
1949 }
1950 }
1951 String::new()
1952}
1953
1954fn stamp_llm(
1959 model: &str,
1960 d: &crate::llm::LlmDraft,
1961 target_ref: String,
1962 cited: Vec<String>,
1963 confidence: f64,
1964 now_ms: i64,
1965) -> Recommendation {
1966 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
1967 let mut args = serde_json::Map::new();
1968 args.insert("text".into(), Value::from(summary_text));
1969 let guidance = if d.guidance.trim().is_empty() {
1970 None
1971 } else {
1972 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
1973 };
1974 let action = ActionKind::Flag;
1975 let mut data = serde_json::Map::new();
1976 data.insert("source".into(), Value::from("llm"));
1977 Recommendation {
1978 hash: String::new(),
1979 analyzer: "loop.llm/1".to_string(),
1980 params_snapshot: serde_json::Map::new(),
1981 origin: Origin::Llm { model: model.to_string() },
1982 target_ref: target_ref.clone(),
1983 action_kind: action,
1984 dedup_key: dedup_key("llm", &target_ref, action),
1985 summary: Summary::new("llm.discover", args),
1986 severity: Severity::Low,
1987 proposal: Proposal::Data { data },
1988 destructive: false,
1989 rollbackable: false,
1990 evidence: cited,
1991 evidence_query: None,
1992 metric: None,
1993 confidence: confidence.clamp(0.0, 1.0),
1995 importance: 0.3,
1996 created_at_ms: now_ms,
1997 guidance,
1998 evalset_hash: None,
1999 status: RecStatus::Pending,
2000 }
2001}
2002
2003fn requires_gating(kind: ActionKind) -> bool {
2008 matches!(
2009 kind,
2010 ActionKind::CodeRevision | ActionKind::AdapterRevision
2011 )
2012}
2013
2014fn stamp(
2015 m: &AnalyzerManifest,
2016 params: &crate::manifest::Params,
2017 d: crate::recommendation::RecDraft,
2018 now_ms: i64,
2019) -> Result<Recommendation> {
2020 let target = TargetRef::parse(&d.target_ref)?;
2021 crate::recommendation::validate_code_rules(
2025 d.action_kind,
2026 target.target_class(),
2027 d.evalset_hash.as_deref(),
2028 )?;
2029 let dedup = dedup_key(m.family(), &d.target_ref, d.action_kind);
2030 let destructive = match &d.proposal {
2031 Proposal::Cal { cal } => cal::contains_destructive(cal),
2032 _ => false,
2033 };
2034 let rollbackable = match &d.proposal {
2035 Proposal::Cal { .. } => !destructive,
2036 Proposal::Edit { .. } => false,
2037 Proposal::Data { .. } => requires_gating(d.action_kind),
2041 };
2042 let mut evidence = d.evidence;
2043 evidence.truncate(MAX_EVIDENCE);
2044 let origin = match m.trust_class {
2049 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
2050 _ => Origin::Builtin,
2051 };
2052 Ok(Recommendation {
2053 hash: String::new(),
2054 analyzer: m.id.clone(),
2055 params_snapshot: params.snapshot(),
2056 origin,
2057 target_ref: target.as_string(),
2058 action_kind: d.action_kind,
2059 dedup_key: dedup,
2060 summary: d.summary,
2061 severity: d.severity,
2062 proposal: d.proposal,
2063 destructive,
2064 rollbackable,
2065 evidence,
2066 evidence_query: d.evidence_query,
2067 metric: d.metric,
2068 confidence: d.confidence,
2069 importance: d.importance,
2070 created_at_ms: now_ms,
2071 guidance: None,
2072 evalset_hash: d.evalset_hash,
2073 status: RecStatus::Pending,
2074 })
2075}
2076
2077fn validate_because(because: &str) -> Result<String> {
2078 let trimmed = because.trim();
2079 if trimmed.is_empty() {
2080 return Err(Error::InvalidProposal(
2081 "a BECAUSE reason is required".into(),
2082 ));
2083 }
2084 if trimmed.chars().count() > MAX_BECAUSE {
2085 return Err(Error::InvalidProposal(format!(
2086 "BECAUSE exceeds {MAX_BECAUSE} chars"
2087 )));
2088 }
2089 Ok(trimmed.to_string())
2090}
2091
2092fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
2093 for req in &m.requires {
2094 match req {
2095 Capability::Forks if !caps.forks => return Some("forks"),
2096 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
2097 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
2098 _ => {}
2099 }
2100 }
2101 None
2102}
2103
2104fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
2105 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
2106}
2107
2108fn gate(
2109 opts: &RunOptions,
2110 p: &LoopPersisted,
2111 new_grains: u64,
2112 new_errors: u64,
2113 now_ms: i64,
2114) -> Option<SkipReason> {
2115 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
2116 if !any {
2117 return None;
2118 }
2119 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
2120 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
2121 let stale_ok = opts
2122 .if_stale_ms
2123 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
2124 if min_new_ok || min_err_ok || stale_ok {
2125 return None;
2126 }
2127 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
2129 Some(SkipReason::NotStale)
2130 } else {
2131 Some(SkipReason::MinNewNotMet)
2132 }
2133}
2134
2135fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
2136 let opts = ReadOpts {
2137 live_only: false,
2138 since_ms: watermark.map(|w| w + 1),
2139 };
2140 let mut new_grains = 0u64;
2141 let mut new_errors = 0u64;
2142 for t in [
2143 crate::model::grain_type::FACT,
2144 crate::model::grain_type::EVENT,
2145 crate::model::grain_type::TOOL,
2146 crate::model::grain_type::OBSERVATION,
2147 ] {
2148 let g = sub.grains_of_type(t, None, opts)?;
2149 new_grains += g.len() as u64;
2150 if t == crate::model::grain_type::TOOL {
2152 new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
2153 }
2154 }
2155 Ok((new_grains, new_errors))
2156}
2157
2158fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
2159 let grains = sub.grains_of_type(
2160 crate::model::grain_type::RECOMMENDATION,
2161 Some(LOOP_NS),
2162 ReadOpts {
2163 live_only: false,
2164 since_ms: None,
2165 },
2166 )?;
2167 let mut set = BTreeSet::new();
2168 for g in grains {
2169 let status = p
2170 .status_index
2171 .get(&g.hash)
2172 .copied()
2173 .unwrap_or(RecStatus::Pending);
2174 if matches!(
2179 status,
2180 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
2181 ) {
2182 if let Some(key) = g.str_field("dedup_key") {
2183 set.insert(key.to_string());
2184 }
2185 }
2186 }
2187 Ok(set)
2188}
2189
2190fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
2191 let g = sub
2192 .grain(rec_hash)?
2193 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
2194 Recommendation::from_fields(rec_hash, &g.fields)
2195}
2196
2197pub(crate) fn is_definition_statement(line: &str) -> bool {
2205 let up = line.trim_start().to_ascii_uppercase();
2206 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
2207}
2208
2209const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
2212 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
2213 approve it to acknowledge it and let it expire.";
2214
2215const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
2220 (evalset hash + run id + stats) — use apply_gated";
2221
2222const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
2223 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
2224 acknowledge it and let it expire.";
2225
2226pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
2239 match proposal {
2240 Proposal::Cal { .. } => Ok(()),
2241 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
2242 Proposal::Data { data } => {
2248 if requires_gating(action_kind)
2249 || data.get("revert_of").and_then(Value::as_str).is_some()
2250 {
2251 Ok(())
2252 } else {
2253 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
2254 }
2255 }
2256 }
2257}
2258
2259#[cfg(test)]
2260mod definition_proposal_tests {
2261 use super::is_definition_statement;
2262
2263 #[test]
2264 fn definition_statements_are_recognized_in_both_spellings() {
2265 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
2266 assert!(is_definition_statement(" define template foo AS { x }"));
2267 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
2268 assert!(!is_definition_statement("ADD fact {}"));
2270 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
2271 assert!(!is_definition_statement("FORGET abc"));
2272 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
2276 }
2277}