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>(
1026 &self,
1027 sub: &S,
1028 rec_hash: &str,
1029 scopes: &ScopeSet,
1030 allow_destructive: bool,
1031 ) -> Result<()> {
1032 if !scopes.has(Scope::Apply) {
1033 return Err(Error::ScopeDenied("apply".into()));
1034 }
1035 let rec = load_rec(sub, rec_hash)?;
1036 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1037 return Err(Error::DestructiveGated(
1038 "destructive apply requires admin scope + allow_destructive".into(),
1039 ));
1040 }
1041 ensure_executable(&rec.proposal)?;
1042 Ok(())
1043 }
1044
1045 #[allow(clippy::too_many_arguments)]
1049 pub fn apply<S: OmsSubstrate>(
1050 &self,
1051 sub: &mut S,
1052 rec_hash: &str,
1053 actor: &str,
1054 observer: ObserverType,
1055 scopes: &ScopeSet,
1056 because: &str,
1057 allow_destructive: bool,
1058 now_ms: i64,
1059 ) -> Result<AppliedRecord> {
1060 self.apply_inner(
1061 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1062 )
1063 }
1064
1065 #[allow(clippy::too_many_arguments)]
1069 pub fn apply_gated<S: OmsSubstrate>(
1070 &self,
1071 sub: &mut S,
1072 rec_hash: &str,
1073 actor: &str,
1074 observer: ObserverType,
1075 scopes: &ScopeSet,
1076 because: &str,
1077 allow_destructive: bool,
1078 gating: &crate::recommendation::GatingEvidence,
1079 now_ms: i64,
1080 ) -> Result<AppliedRecord> {
1081 self.apply_inner(
1082 sub,
1083 rec_hash,
1084 actor,
1085 observer,
1086 scopes,
1087 because,
1088 allow_destructive,
1089 Some(gating),
1090 now_ms,
1091 )
1092 }
1093
1094 #[allow(clippy::too_many_arguments)]
1095 fn apply_inner<S: OmsSubstrate>(
1096 &self,
1097 sub: &mut S,
1098 rec_hash: &str,
1099 actor: &str,
1100 observer: ObserverType,
1101 scopes: &ScopeSet,
1102 because: &str,
1103 allow_destructive: bool,
1104 gating: Option<&crate::recommendation::GatingEvidence>,
1105 now_ms: i64,
1106 ) -> Result<AppliedRecord> {
1107 if !scopes.has(Scope::Apply) {
1108 return Err(Error::ScopeDenied("apply".into()));
1109 }
1110 let because = validate_because(because)?;
1111 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1112 let status = *p
1113 .status_index
1114 .get(rec_hash)
1115 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1116 if !status.can_transition_to(RecStatus::Applied, false) {
1117 return Err(Error::LifecycleViolation(format!(
1118 "{} -> applied (approve first)",
1119 status.as_str()
1120 )));
1121 }
1122 let rec = load_rec(sub, rec_hash)?;
1123 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1124 return Err(Error::DestructiveGated(
1125 "destructive apply requires admin scope + allow_destructive".into(),
1126 ));
1127 }
1128 if rec.action_kind == ActionKind::CodeRevision {
1133 let g = gating.ok_or_else(|| {
1134 Error::InvalidProposal(
1135 "code revisions apply only with a recorded gating run \
1136 (evalset hash + run id + stats) — use apply_gated"
1137 .into(),
1138 )
1139 })?;
1140 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1141 if g.evalset_hash != pin {
1142 return Err(Error::InvalidProposal(format!(
1143 "gating ran evalset {} but the recommendation is pinned \
1144 to {pin} (Rule E1)",
1145 g.evalset_hash
1146 )));
1147 }
1148 match sub.grain(pin)? {
1149 Some(evalset) if evalset.is_live() => {}
1150 Some(_) => {
1151 return Err(Error::InvalidProposal(
1152 "the pinned evalset was superseded after gating — \
1153 the recommendation must re-gate (Rule E1)"
1154 .into(),
1155 ))
1156 }
1157 None => {
1158 return Err(Error::InvalidProposal(format!(
1159 "pinned evalset {pin} not found in the substrate"
1160 )))
1161 }
1162 }
1163 if g.failed > 0 {
1164 return Err(Error::InvalidProposal(format!(
1165 "the gating run failed {}/{} cases — a failing gate \
1166 cannot admit code",
1167 g.failed,
1168 g.passed + g.failed
1169 )));
1170 }
1171 }
1172
1173 let mut created = Vec::new();
1175 let mut inverse_cal: Option<String> = None;
1179 match &rec.proposal {
1180 Proposal::Cal { cal } => {
1181 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1182 if !is_definition_statement(line) {
1183 continue;
1184 }
1185 match sub.definition_inverse(line)? {
1186 Some(inv) => inverse_cal = Some(inv),
1187 None => {
1188 return Err(Error::InvalidProposal(format!(
1189 "this substrate cannot record a rollback inverse for {line:?}; \
1190 a definition rewrite that ROLLBACK could not undo is refused \
1191 rather than applied"
1192 )))
1193 }
1194 }
1195 }
1196 let rows = sub.execute_cal(cal)?;
1197 for r in rows {
1198 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1199 created.push(h.to_string());
1200 }
1201 }
1202 }
1203 Proposal::Edit { .. } => {
1206 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1207 }
1208 Proposal::Data { data } if rec.action_kind == ActionKind::CodeRevision => {
1215 let mut spec = crate::substrate::GrainSpec::new(
1216 crate::model::grain_type::FACT,
1217 LOOP_NS,
1218 )
1219 .with_field("subject", rec.target_ref.clone())
1220 .with_field("relation", "mg:code_promotion")
1221 .with_field(
1222 "object",
1223 serde_json::to_string(data)
1224 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1225 )
1226 .with_field("rec_hash", rec_hash.to_string());
1227 if let Some(g) = gating {
1228 spec = spec
1229 .with_field("gating_evalset", g.evalset_hash.clone())
1230 .with_field("gating_run_id", g.run_id.clone());
1231 }
1232 created.push(sub.put_grain(&spec)?);
1233 }
1234 Proposal::Data { data } => {
1235 let revert_of = data
1240 .get("revert_of")
1241 .and_then(Value::as_str)
1242 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1243 self.rollback(
1244 sub,
1245 revert_of,
1246 actor,
1247 observer,
1248 scopes,
1249 &because,
1250 now_ms,
1251 )?;
1252 p = LoopPersisted::from_value(sub.load_state()?)?;
1255 }
1256 }
1257
1258 let applied = AppliedRecord {
1259 applied_at_ms: now_ms,
1260 target_ref: rec.target_ref.clone(),
1261 rollbackable: rec.rollbackable,
1262 created_hashes: created,
1263 inverse_cal,
1264 metric: rec.metric.clone(),
1265 };
1266 let prev = p.audit_heads.get(rec_hash).cloned();
1267 let audit = AuditRecord {
1268 rec_hash: rec_hash.into(),
1269 from: Some(status),
1270 to: RecStatus::Applied,
1271 actor: actor.into(),
1272 observer_type: observer,
1273 because,
1274 previous_audit_hash: prev,
1275 gating: gating.cloned(),
1276 at_ms: now_ms,
1277 };
1278 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1279 p.audit_heads.insert(rec_hash.into(), audit_hash);
1280 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1281 p.applied.insert(rec_hash.into(), applied.clone());
1282 sub.store_state(&p.to_value()?)?;
1283 Ok(applied)
1284 }
1285
1286 #[allow(clippy::too_many_arguments)]
1289 pub fn rollback<S: OmsSubstrate>(
1290 &self,
1291 sub: &mut S,
1292 rec_hash: &str,
1293 actor: &str,
1294 observer: ObserverType,
1295 scopes: &ScopeSet,
1296 because: &str,
1297 now_ms: i64,
1298 ) -> Result<()> {
1299 if !scopes.has(Scope::Apply) {
1300 return Err(Error::ScopeDenied("apply".into()));
1301 }
1302 let because = validate_because(because)?;
1303 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1304 let status = *p
1305 .status_index
1306 .get(rec_hash)
1307 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1308 if !status.can_transition_to(RecStatus::RolledBack, false) {
1309 return Err(Error::LifecycleViolation(format!(
1310 "{} -> rolled_back",
1311 status.as_str()
1312 )));
1313 }
1314 let applied = p
1315 .applied
1316 .get(rec_hash)
1317 .cloned()
1318 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1319 if !applied.rollbackable {
1320 return Err(Error::LifecycleViolation(
1321 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1322 ));
1323 }
1324 for h in &applied.created_hashes {
1325 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1326 }
1327 if let Some(inverse) = &applied.inverse_cal {
1334 sub.execute_cal(inverse)?;
1335 }
1336 let prev = p.audit_heads.get(rec_hash).cloned();
1337 let audit = AuditRecord {
1338 rec_hash: rec_hash.into(),
1339 from: Some(status),
1340 to: RecStatus::RolledBack,
1341 actor: actor.into(),
1342 observer_type: observer,
1343 because,
1344 previous_audit_hash: prev,
1345 gating: None,
1346 at_ms: now_ms,
1347 };
1348 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1349 p.audit_heads.insert(rec_hash.into(), audit_hash);
1350 p.status_index
1351 .insert(rec_hash.into(), RecStatus::RolledBack);
1352 sub.store_state(&p.to_value()?)?;
1353 Ok(())
1354 }
1355
1356 pub fn recommendations<S: OmsSubstrate>(
1361 &self,
1362 sub: &S,
1363 status_filter: Option<RecStatus>,
1364 ) -> Result<Vec<Recommendation>> {
1365 let p = LoopPersisted::from_value(sub.load_state()?)?;
1366 let grains = sub.grains_of_type(
1367 crate::model::grain_type::RECOMMENDATION,
1368 Some(LOOP_NS),
1369 ReadOpts {
1370 live_only: false,
1371 since_ms: None,
1372 },
1373 )?;
1374 let mut out = Vec::new();
1375 for g in grains {
1376 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1377 rec.status = p
1378 .status_index
1379 .get(&g.hash)
1380 .copied()
1381 .unwrap_or(RecStatus::Pending);
1382 if let Some(f) = status_filter {
1383 if rec.status != f {
1384 continue;
1385 }
1386 }
1387 out.push(rec);
1388 }
1389 out.sort_by(|a, b| {
1398 b.severity
1399 .cmp(&a.severity)
1400 .then(a.created_at_ms.cmp(&b.created_at_ms))
1401 .then(a.dedup_key.cmp(&b.dedup_key))
1402 .then(a.hash.cmp(&b.hash))
1403 });
1404 Ok(out)
1405 }
1406
1407 pub fn analyzer_settings<S: OmsSubstrate>(
1410 &self,
1411 sub: &S,
1412 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1413 let p = LoopPersisted::from_value(sub.load_state()?)?;
1414 Ok(self
1415 .analyzers
1416 .iter()
1417 .map(|a| {
1418 let m = a.manifest();
1419 let cfg = p.config.get(&m.id);
1420 crate::config::AnalyzerSetting {
1421 id: m.id.clone(),
1422 title: m.title.clone(),
1423 description: m.description.clone(),
1424 tier: format!("{:?}", m.tier),
1425 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1426 default_on: m.default_on,
1427 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1428 severity_floor: cfg
1429 .and_then(|c| c.severity_floor)
1430 .map(|s| s.as_str().to_string()),
1431 }
1432 })
1433 .collect())
1434 }
1435
1436 pub fn set_analyzer_config<S: OmsSubstrate>(
1442 &self,
1443 sub: &mut S,
1444 analyzer_id: &str,
1445 update: crate::config::AnalyzerConfigUpdate,
1446 scopes: &ScopeSet,
1447 ) -> Result<crate::config::AnalyzerConfig> {
1448 if !scopes.has(Scope::Admin) {
1449 return Err(Error::ScopeDenied("admin".into()));
1450 }
1451 let manifest = self
1452 .analyzers
1453 .iter()
1454 .map(|a| a.manifest())
1455 .find(|m| m.id == analyzer_id)
1456 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1457 if let Some(params) = &update.params {
1459 manifest.resolve_params(params)?;
1460 }
1461 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1462 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1463 if let Some(enabled) = update.enabled {
1464 cfg.enabled = Some(enabled);
1465 }
1466 if update.clear_floor {
1467 cfg.severity_floor = None;
1468 } else if let Some(floor) = update.severity_floor {
1469 cfg.severity_floor = Some(floor);
1470 }
1471 if let Some(params) = update.params {
1472 cfg.params = params;
1473 }
1474 if let Some(ns) = update.namespaces {
1475 cfg.namespaces = ns;
1476 }
1477 let stored = cfg.clone();
1478 sub.store_state(&p.to_value()?)?;
1479 Ok(stored)
1480 }
1481
1482 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1485 let p = LoopPersisted::from_value(sub.load_state()?)?;
1486 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1487 out.sort_by(|a, b| {
1491 a.measured_at_ms
1492 .cmp(&b.measured_at_ms)
1493 .then(a.horizon_ms.cmp(&b.horizon_ms))
1494 .then(a.metric.cmp(&b.metric))
1495 .then(a.rec_hash.cmp(&b.rec_hash))
1496 });
1497 Ok(out)
1498 }
1499
1500 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1504 let p = LoopPersisted::from_value(sub.load_state()?)?;
1505 let (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
1506 let recs = self.recommendations(sub, None)?;
1507 let mut pending = 0;
1508 let mut applied = 0;
1509 for r in &recs {
1510 match r.status {
1511 RecStatus::Pending => pending += 1,
1512 RecStatus::Applied => applied += 1,
1513 _ => {}
1514 }
1515 }
1516 let stale = match p.state.last_run_ms {
1518 None => true,
1519 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1520 };
1521 Ok(Health {
1522 last_run_ms: p.state.last_run_ms,
1523 grains_since_run,
1524 error_events_since_run,
1525 pending,
1526 applied,
1527 total: recs.len() as u64,
1528 stale,
1529 })
1530 }
1531
1532 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1537 let recs = self.recommendations(sub, None)?;
1538 let mut m = LlmMetrics::default();
1539 for r in &recs {
1540 if !matches!(r.origin, Origin::Llm { .. }) {
1541 continue;
1542 }
1543 m.proposed += 1;
1544 match r.status {
1545 RecStatus::Pending => m.pending += 1,
1546 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1547 RecStatus::Rejected => m.rejected += 1,
1548 RecStatus::Expired => {}
1549 }
1550 }
1551 let decided = m.approved + m.rejected;
1552 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1553 Ok(m)
1554 }
1555}
1556
1557#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1559pub struct Health {
1560 #[serde(skip_serializing_if = "Option::is_none")]
1561 pub last_run_ms: Option<i64>,
1562 pub grains_since_run: u64,
1563 pub error_events_since_run: u64,
1564 pub pending: u64,
1565 pub applied: u64,
1566 pub total: u64,
1567 pub stale: bool,
1570}
1571
1572#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1574pub struct LlmMetrics {
1575 pub proposed: u64,
1578 pub pending: u64,
1579 pub approved: u64,
1581 pub rejected: u64,
1582 #[serde(skip_serializing_if = "Option::is_none")]
1584 pub approval_rate: Option<f64>,
1585}
1586
1587fn measure_outcomes<S: OmsSubstrate>(
1594 sub: &S,
1595 p: &mut LoopPersisted,
1596 now_ms: i64,
1597) -> Result<Vec<OutcomeInput>> {
1598 let mut due: Vec<(String, crate::config::AppliedRecord, i64)> = Vec::new();
1600 for (h, a) in &p.applied {
1601 if p.status_index.get(h) != Some(&RecStatus::Applied) {
1602 continue;
1603 }
1604 let Some(metric) = &a.metric else { continue };
1605 let done = p.measured.get(h).cloned().unwrap_or_default();
1606 for horizon in metric.horizons() {
1607 if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
1608 due.push((h.clone(), a.clone(), horizon));
1609 }
1610 }
1611 }
1612
1613 let mut out = Vec::new();
1614 for (rec_hash, applied, horizon) in due {
1615 let metric = applied.metric.as_ref().unwrap();
1616 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1617 continue; };
1619 let regressed = crate::recommendation::is_regression(
1620 metric.baseline,
1621 current,
1622 metric.higher_is_better,
1623 );
1624 p.outcomes.entry(rec_hash.clone()).or_default().push(
1625 crate::recommendation::OutcomeResult {
1626 rec_hash: rec_hash.clone(),
1627 metric: metric.metric.clone(),
1628 baseline: metric.baseline,
1629 current,
1630 verdict: if regressed { "regressed" } else { "held" }.into(),
1631 horizon_ms: horizon,
1632 measured_at_ms: now_ms,
1633 },
1634 );
1635 p.measured.entry(rec_hash.clone()).or_default().push(horizon);
1636 if regressed {
1637 out.push(OutcomeInput {
1638 rec_hash,
1639 target_ref: applied.target_ref.clone(),
1640 metric: metric.metric.clone(),
1641 baseline: metric.baseline,
1642 current,
1643 unit: metric.unit.clone(),
1644 higher_is_better: metric.higher_is_better,
1645 });
1646 }
1647 }
1648 Ok(out)
1649}
1650
1651pub(crate) fn measure_metric<S: SubstrateRead>(
1653 sub: &S,
1654 metric: &crate::recommendation::MetricSnapshot,
1655 since_ms: i64,
1656) -> Result<Option<f64>> {
1657 match metric.metric.as_str() {
1658 "tool_error_recurrence" => {
1663 let Some(tool) = &metric.subject else { return Ok(None) };
1664 let tools = sub.grains_of_type(
1665 crate::model::grain_type::TOOL,
1666 None,
1667 ReadOpts { live_only: true, since_ms: Some(since_ms) },
1668 )?;
1669 let n = tools
1670 .iter()
1671 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
1672 .filter(|t| {
1673 metric.relation.as_deref().is_none_or(|sig| {
1676 crate::analyzers::tool_failure::normalize_signature(
1677 t.tool_content().unwrap_or(""),
1678 ) == sig
1679 })
1680 })
1681 .count();
1682 Ok(Some(n as f64))
1683 }
1684 "contradiction_recurrence" => {
1688 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
1689 return Ok(None);
1690 };
1691 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
1692 let distinct: BTreeSet<String> = facts
1693 .iter()
1694 .filter(|f| {
1695 f.fact_relation()
1696 .is_some_and(|r| normalize_ident(r) == *relation)
1697 })
1698 .filter_map(|f| f.fact_object().map(normalize_ident))
1699 .collect();
1700 Ok(Some(distinct.len().saturating_sub(1) as f64))
1701 }
1702 m if m.starts_with("evalset:") => {
1715 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
1716 return Ok(None);
1717 };
1718 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
1719 return Ok(None);
1720 };
1721 Ok(match field {
1725 "failed" => Some(run.failed as f64),
1726 "passed" => Some(run.passed as f64),
1727 "total" => Some(run.total() as f64),
1728 "error_rate" => match run.total() {
1729 0 => None, t => Some(run.failed as f64 / t as f64),
1731 },
1732 other => run.field(other),
1733 })
1734 }
1735 _ => Ok(None),
1736 }
1737}
1738
1739fn scoped_live_facts<S: SubstrateRead>(
1742 sub: &S,
1743 namespace: Option<&str>,
1744 subject: &str,
1745) -> Result<Vec<GrainRecord>> {
1746 let facts = sub.grains_of_type(
1747 crate::model::grain_type::FACT,
1748 None,
1749 ReadOpts { live_only: true, since_ms: None },
1750 )?;
1751 Ok(facts
1752 .into_iter()
1753 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
1754 .filter(|f| {
1755 f.fact_subject()
1756 .is_some_and(|s| normalize_ident(s) == subject)
1757 })
1758 .collect())
1759}
1760
1761fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
1773 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
1774 return false;
1775 };
1776 if fields.is_empty() {
1777 return false;
1778 }
1779 let Ok(Some(grain)) = sub.grain(&target) else {
1780 return false;
1781 };
1782 if grain.valid_to_ms.is_some() {
1792 return false;
1793 }
1794 fields.iter().all(|(k, v)| {
1796 if k == "namespace" {
1797 return v
1798 .as_str()
1799 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
1800 }
1801 match (v, grain.fields.get(k)) {
1802 (Value::String(a), Some(Value::String(b))) => {
1803 normalize_ident(a) == normalize_ident(b)
1804 }
1805 (a, Some(b)) => a == b,
1806 (_, None) => false,
1807 }
1808 })
1809}
1810
1811const MIN_LLM_CONFIDENCE: f64 = 0.75;
1814
1815const DISCOVER_INSTRUCTIONS: &str = "You review an agent's memory for quality. \
1820Given deterministic findings and the evidence they cite, propose ADDITIONAL \
1821findings the deterministic checks would miss (e.g. a semantic contradiction, a \
1822stale assumption, a duplicated meaning). SCORING: propose a finding ONLY if you \
1823are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
1824useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
1825earns 0. When in doubt, propose nothing — an empty list is the correct answer \
1826when there is nothing worth flagging. The 'approved' and 'rejected' lists, when \
1827present, show findings this reviewer recently accepted or rejected — prefer the \
1828kind they accept and avoid the kind they reject. Every proposal MUST cite one \
1829or more evidence hashes from the bundle, target a memory entity, and include \
1830your confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
1831\"target\":\"entity:<ns>/<subject>\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
1832\"confidence\":0.0}]}. Propose nothing you cannot ground in the evidence.";
1833
1834const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
1840fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
1841the facts it relies on are actually present in the cited evidence, NOT that its \
1842conclusion is stated verbatim. Decompose the finding into the factual claims it \
1843depends on. Mark supported=true when those facts are present in the evidence \
1844(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
1845on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
1846different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
1847
1848const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
1850each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
1851never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
1852SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
1853'possible' findings with no concrete defect, and reject any claimed \
1854inconsistency or contradiction that is not backed by at least two actually \
1855conflicting facts in the cited evidence. (2) Context — does the finding \
1856correctly read its cited evidence, or misinterpret what the grains say? Keep a \
1857finding when it names a genuine, specific problem grounded in its evidence and \
1858materially useful to a human reviewer; otherwise reject it, and default to \
1859keep=false when uncertain. Do NOT reject a finding for being 'already known', \
1860redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
1861grounded cross-fact inconsistency with two conflicting facts is exactly what to \
1862KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
1863{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
1864
1865const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
1867guidance note to help a human reviewer decide. Do not restate the finding. Return \
1868JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
1869
1870fn push_evidence(
1873 evidence: &mut Vec<crate::llm::EvidenceItem>,
1874 bundle: &mut BTreeSet<String>,
1875 g: &GrainRecord,
1876) {
1877 if evidence.len() < 64 && bundle.insert(g.hash.clone()) {
1878 evidence.push(crate::llm::EvidenceItem {
1879 hash: g.hash.clone(),
1880 grain_type: g.grain_type.clone(),
1881 text: crate::llm::cap(&grain_brief(g), 400),
1882 });
1883 }
1884}
1885
1886fn grain_brief(g: &GrainRecord) -> String {
1888 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
1889 return format!("{s} {r} {o}");
1890 }
1891 for key in ["content", "body", "text", "summary"] {
1892 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
1893 if !v.is_empty() {
1894 return v.to_string();
1895 }
1896 }
1897 }
1898 String::new()
1899}
1900
1901fn stamp_llm(
1906 model: &str,
1907 d: &crate::llm::LlmDraft,
1908 target_ref: String,
1909 cited: Vec<String>,
1910 confidence: f64,
1911 now_ms: i64,
1912) -> Recommendation {
1913 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
1914 let mut args = serde_json::Map::new();
1915 args.insert("text".into(), Value::from(summary_text));
1916 let guidance = if d.guidance.trim().is_empty() {
1917 None
1918 } else {
1919 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
1920 };
1921 let action = ActionKind::Flag;
1922 let mut data = serde_json::Map::new();
1923 data.insert("source".into(), Value::from("llm"));
1924 Recommendation {
1925 hash: String::new(),
1926 analyzer: "loop.llm/1".to_string(),
1927 params_snapshot: serde_json::Map::new(),
1928 origin: Origin::Llm { model: model.to_string() },
1929 target_ref: target_ref.clone(),
1930 action_kind: action,
1931 dedup_key: dedup_key("llm", &target_ref, action),
1932 summary: Summary::new("llm.discover", args),
1933 severity: Severity::Low,
1934 proposal: Proposal::Data { data },
1935 destructive: false,
1936 rollbackable: false,
1937 evidence: cited,
1938 evidence_query: None,
1939 metric: None,
1940 confidence: confidence.clamp(0.0, 1.0),
1942 importance: 0.3,
1943 created_at_ms: now_ms,
1944 guidance,
1945 evalset_hash: None,
1946 status: RecStatus::Pending,
1947 }
1948}
1949
1950fn stamp(
1951 m: &AnalyzerManifest,
1952 params: &crate::manifest::Params,
1953 d: crate::recommendation::RecDraft,
1954 now_ms: i64,
1955) -> Result<Recommendation> {
1956 let target = TargetRef::parse(&d.target_ref)?;
1957 crate::recommendation::validate_code_rules(
1961 d.action_kind,
1962 target.target_class(),
1963 d.evalset_hash.as_deref(),
1964 )?;
1965 let dedup = dedup_key(m.family(), &d.target_ref, d.action_kind);
1966 let destructive = match &d.proposal {
1967 Proposal::Cal { cal } => cal::contains_destructive(cal),
1968 _ => false,
1969 };
1970 let rollbackable = match &d.proposal {
1971 Proposal::Cal { .. } => !destructive,
1972 Proposal::Edit { .. } => false,
1973 Proposal::Data { .. } => d.action_kind == ActionKind::CodeRevision,
1976 };
1977 let mut evidence = d.evidence;
1978 evidence.truncate(MAX_EVIDENCE);
1979 let origin = match m.trust_class {
1984 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
1985 _ => Origin::Builtin,
1986 };
1987 Ok(Recommendation {
1988 hash: String::new(),
1989 analyzer: m.id.clone(),
1990 params_snapshot: params.snapshot(),
1991 origin,
1992 target_ref: target.as_string(),
1993 action_kind: d.action_kind,
1994 dedup_key: dedup,
1995 summary: d.summary,
1996 severity: d.severity,
1997 proposal: d.proposal,
1998 destructive,
1999 rollbackable,
2000 evidence,
2001 evidence_query: d.evidence_query,
2002 metric: d.metric,
2003 confidence: d.confidence,
2004 importance: d.importance,
2005 created_at_ms: now_ms,
2006 guidance: None,
2007 evalset_hash: d.evalset_hash,
2008 status: RecStatus::Pending,
2009 })
2010}
2011
2012fn validate_because(because: &str) -> Result<String> {
2013 let trimmed = because.trim();
2014 if trimmed.is_empty() {
2015 return Err(Error::InvalidProposal(
2016 "a BECAUSE reason is required".into(),
2017 ));
2018 }
2019 if trimmed.chars().count() > MAX_BECAUSE {
2020 return Err(Error::InvalidProposal(format!(
2021 "BECAUSE exceeds {MAX_BECAUSE} chars"
2022 )));
2023 }
2024 Ok(trimmed.to_string())
2025}
2026
2027fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
2028 for req in &m.requires {
2029 match req {
2030 Capability::Forks if !caps.forks => return Some("forks"),
2031 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
2032 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
2033 _ => {}
2034 }
2035 }
2036 None
2037}
2038
2039fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
2040 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
2041}
2042
2043fn gate(
2044 opts: &RunOptions,
2045 p: &LoopPersisted,
2046 new_grains: u64,
2047 new_errors: u64,
2048 now_ms: i64,
2049) -> Option<SkipReason> {
2050 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
2051 if !any {
2052 return None;
2053 }
2054 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
2055 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
2056 let stale_ok = opts
2057 .if_stale_ms
2058 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
2059 if min_new_ok || min_err_ok || stale_ok {
2060 return None;
2061 }
2062 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
2064 Some(SkipReason::NotStale)
2065 } else {
2066 Some(SkipReason::MinNewNotMet)
2067 }
2068}
2069
2070fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
2071 let opts = ReadOpts {
2072 live_only: false,
2073 since_ms: watermark.map(|w| w + 1),
2074 };
2075 let mut new_grains = 0u64;
2076 let mut new_errors = 0u64;
2077 for t in [
2078 crate::model::grain_type::FACT,
2079 crate::model::grain_type::EVENT,
2080 crate::model::grain_type::TOOL,
2081 crate::model::grain_type::OBSERVATION,
2082 ] {
2083 let g = sub.grains_of_type(t, None, opts)?;
2084 new_grains += g.len() as u64;
2085 if t == crate::model::grain_type::TOOL {
2087 new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
2088 }
2089 }
2090 Ok((new_grains, new_errors))
2091}
2092
2093fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
2094 let grains = sub.grains_of_type(
2095 crate::model::grain_type::RECOMMENDATION,
2096 Some(LOOP_NS),
2097 ReadOpts {
2098 live_only: false,
2099 since_ms: None,
2100 },
2101 )?;
2102 let mut set = BTreeSet::new();
2103 for g in grains {
2104 let status = p
2105 .status_index
2106 .get(&g.hash)
2107 .copied()
2108 .unwrap_or(RecStatus::Pending);
2109 if matches!(
2114 status,
2115 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
2116 ) {
2117 if let Some(key) = g.str_field("dedup_key") {
2118 set.insert(key.to_string());
2119 }
2120 }
2121 }
2122 Ok(set)
2123}
2124
2125fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
2126 let g = sub
2127 .grain(rec_hash)?
2128 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
2129 Recommendation::from_fields(rec_hash, &g.fields)
2130}
2131
2132pub(crate) fn is_definition_statement(line: &str) -> bool {
2140 let up = line.trim_start().to_ascii_uppercase();
2141 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
2142}
2143
2144const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
2147 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
2148 approve it to acknowledge it and let it expire.";
2149
2150const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
2153 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
2154 acknowledge it and let it expire.";
2155
2156pub(crate) fn ensure_executable(proposal: &Proposal) -> Result<()> {
2169 match proposal {
2170 Proposal::Cal { .. } => Ok(()),
2171 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
2172 Proposal::Data { data } => {
2175 if data.get("revert_of").and_then(Value::as_str).is_some() {
2176 Ok(())
2177 } else {
2178 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
2179 }
2180 }
2181 }
2182}
2183
2184#[cfg(test)]
2185mod definition_proposal_tests {
2186 use super::is_definition_statement;
2187
2188 #[test]
2189 fn definition_statements_are_recognized_in_both_spellings() {
2190 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
2191 assert!(is_definition_statement(" define template foo AS { x }"));
2192 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
2193 assert!(!is_definition_statement("ADD fact {}"));
2195 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
2196 assert!(!is_definition_statement("FORGET abc"));
2197 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
2201 }
2202}