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 for r in sub.execute_cal(cal)? {
886 if let Some(h) = r.get("hash").and_then(Value::as_str) {
887 created.push(h.to_string());
888 }
889 }
890 }
891 let applied = AppliedRecord {
892 applied_at_ms: now_ms,
893 target_ref: rec.target_ref.clone(),
894 rollbackable: rec.rollbackable,
895 created_hashes: created,
896 metric: rec.metric.clone(),
897 };
898 let prev = p.audit_heads.get(&rec.hash).cloned();
899 let audit = AuditRecord {
900 rec_hash: rec.hash.clone(),
901 from: Some(RecStatus::Pending),
902 to: RecStatus::Applied,
903 actor: "policy:auto".into(),
904 observer_type: ObserverType::Policy,
905 because: "auto-applied per host policy".into(),
906 previous_audit_hash: prev,
907 gating: None,
908 at_ms: now_ms,
909 };
910 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
911 p.audit_heads.insert(rec.hash.clone(), audit_hash);
912 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
913 p.applied.insert(rec.hash.clone(), applied);
914 Ok(())
915 }
916
917 #[allow(clippy::too_many_arguments)]
920 pub fn review<S: OmsSubstrate>(
921 &self,
922 sub: &mut S,
923 rec_hash: &str,
924 decision: Decision,
925 actor: &str,
926 observer: ObserverType,
927 scopes: &ScopeSet,
928 because: &str,
929 now_ms: i64,
930 ) -> Result<()> {
931 if !scopes.has(Scope::Review) {
932 return Err(Error::ScopeDenied("review".into()));
933 }
934 let because = validate_because(because)?;
935 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
936 let status = *p
937 .status_index
938 .get(rec_hash)
939 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
940 let to = match decision {
941 Decision::Approve => RecStatus::Approved,
942 Decision::Reject => RecStatus::Rejected,
943 };
944 if !status.can_transition_to(to, false) {
945 return Err(Error::LifecycleViolation(format!(
946 "{} -> {}",
947 status.as_str(),
948 to.as_str()
949 )));
950 }
951 if to == RecStatus::Approved {
952 if let Some(creator) = p.creators.get(rec_hash) {
953 if creator == actor {
954 return Err(Error::SelfApproval(format!(
955 "{actor} created this recommendation"
956 )));
957 }
958 }
959 if let Some(trigger) = p.co_creators.get(rec_hash) {
960 if trigger == actor {
961 return Err(Error::SelfApproval(format!(
962 "{actor} triggered the run that authored this recommendation"
963 )));
964 }
965 }
966 }
967 let prev = p.audit_heads.get(rec_hash).cloned();
968 let audit = AuditRecord {
969 rec_hash: rec_hash.into(),
970 from: Some(status),
971 to,
972 actor: actor.into(),
973 observer_type: observer,
974 because,
975 previous_audit_hash: prev,
976 gating: None,
977 at_ms: now_ms,
978 };
979 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
980 p.audit_heads.insert(rec_hash.into(), audit_hash);
981 p.status_index.insert(rec_hash.into(), to);
982 if to == RecStatus::Rejected {
983 if let Ok(rec) = load_rec(sub, rec_hash) {
984 const BASE_MS: i64 = 7 * 86_400_000;
989 const CAP_MS: i64 = 90 * 86_400_000;
990 let strikes = p.cooldown_strikes.entry(rec.dedup_key.clone()).or_insert(0);
991 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
992 *strikes = strikes.saturating_add(1);
993 p.cooldowns.insert(rec.dedup_key, now_ms + interval);
994 }
995 }
996 sub.store_state(&p.to_value()?)?;
997 Ok(())
998 }
999
1000 pub fn preflight_apply<S: OmsSubstrate>(
1013 &self,
1014 sub: &S,
1015 rec_hash: &str,
1016 scopes: &ScopeSet,
1017 allow_destructive: bool,
1018 ) -> Result<()> {
1019 if !scopes.has(Scope::Apply) {
1020 return Err(Error::ScopeDenied("apply".into()));
1021 }
1022 let rec = load_rec(sub, rec_hash)?;
1023 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1024 return Err(Error::DestructiveGated(
1025 "destructive apply requires admin scope + allow_destructive".into(),
1026 ));
1027 }
1028 ensure_executable(&rec.proposal)?;
1029 Ok(())
1030 }
1031
1032 #[allow(clippy::too_many_arguments)]
1036 pub fn apply<S: OmsSubstrate>(
1037 &self,
1038 sub: &mut S,
1039 rec_hash: &str,
1040 actor: &str,
1041 observer: ObserverType,
1042 scopes: &ScopeSet,
1043 because: &str,
1044 allow_destructive: bool,
1045 now_ms: i64,
1046 ) -> Result<AppliedRecord> {
1047 self.apply_inner(
1048 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1049 )
1050 }
1051
1052 #[allow(clippy::too_many_arguments)]
1056 pub fn apply_gated<S: OmsSubstrate>(
1057 &self,
1058 sub: &mut S,
1059 rec_hash: &str,
1060 actor: &str,
1061 observer: ObserverType,
1062 scopes: &ScopeSet,
1063 because: &str,
1064 allow_destructive: bool,
1065 gating: &crate::recommendation::GatingEvidence,
1066 now_ms: i64,
1067 ) -> Result<AppliedRecord> {
1068 self.apply_inner(
1069 sub,
1070 rec_hash,
1071 actor,
1072 observer,
1073 scopes,
1074 because,
1075 allow_destructive,
1076 Some(gating),
1077 now_ms,
1078 )
1079 }
1080
1081 #[allow(clippy::too_many_arguments)]
1082 fn apply_inner<S: OmsSubstrate>(
1083 &self,
1084 sub: &mut S,
1085 rec_hash: &str,
1086 actor: &str,
1087 observer: ObserverType,
1088 scopes: &ScopeSet,
1089 because: &str,
1090 allow_destructive: bool,
1091 gating: Option<&crate::recommendation::GatingEvidence>,
1092 now_ms: i64,
1093 ) -> Result<AppliedRecord> {
1094 if !scopes.has(Scope::Apply) {
1095 return Err(Error::ScopeDenied("apply".into()));
1096 }
1097 let because = validate_because(because)?;
1098 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1099 let status = *p
1100 .status_index
1101 .get(rec_hash)
1102 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1103 if !status.can_transition_to(RecStatus::Applied, false) {
1104 return Err(Error::LifecycleViolation(format!(
1105 "{} -> applied (approve first)",
1106 status.as_str()
1107 )));
1108 }
1109 let rec = load_rec(sub, rec_hash)?;
1110 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1111 return Err(Error::DestructiveGated(
1112 "destructive apply requires admin scope + allow_destructive".into(),
1113 ));
1114 }
1115 if rec.action_kind == ActionKind::CodeRevision {
1120 let g = gating.ok_or_else(|| {
1121 Error::InvalidProposal(
1122 "code revisions apply only with a recorded gating run \
1123 (evalset hash + run id + stats) — use apply_gated"
1124 .into(),
1125 )
1126 })?;
1127 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1128 if g.evalset_hash != pin {
1129 return Err(Error::InvalidProposal(format!(
1130 "gating ran evalset {} but the recommendation is pinned \
1131 to {pin} (Rule E1)",
1132 g.evalset_hash
1133 )));
1134 }
1135 match sub.grain(pin)? {
1136 Some(evalset) if evalset.is_live() => {}
1137 Some(_) => {
1138 return Err(Error::InvalidProposal(
1139 "the pinned evalset was superseded after gating — \
1140 the recommendation must re-gate (Rule E1)"
1141 .into(),
1142 ))
1143 }
1144 None => {
1145 return Err(Error::InvalidProposal(format!(
1146 "pinned evalset {pin} not found in the substrate"
1147 )))
1148 }
1149 }
1150 if g.failed > 0 {
1151 return Err(Error::InvalidProposal(format!(
1152 "the gating run failed {}/{} cases — a failing gate \
1153 cannot admit code",
1154 g.failed,
1155 g.passed + g.failed
1156 )));
1157 }
1158 }
1159
1160 let mut created = Vec::new();
1162 match &rec.proposal {
1163 Proposal::Cal { cal } => {
1164 let rows = sub.execute_cal(cal)?;
1165 for r in rows {
1166 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1167 created.push(h.to_string());
1168 }
1169 }
1170 }
1171 Proposal::Edit { .. } => {
1174 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1175 }
1176 Proposal::Data { data } if rec.action_kind == ActionKind::CodeRevision => {
1183 let mut spec = crate::substrate::GrainSpec::new(
1184 crate::model::grain_type::FACT,
1185 LOOP_NS,
1186 )
1187 .with_field("subject", rec.target_ref.clone())
1188 .with_field("relation", "mg:code_promotion")
1189 .with_field(
1190 "object",
1191 serde_json::to_string(data)
1192 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1193 )
1194 .with_field("rec_hash", rec_hash.to_string());
1195 if let Some(g) = gating {
1196 spec = spec
1197 .with_field("gating_evalset", g.evalset_hash.clone())
1198 .with_field("gating_run_id", g.run_id.clone());
1199 }
1200 created.push(sub.put_grain(&spec)?);
1201 }
1202 Proposal::Data { data } => {
1203 let revert_of = data
1208 .get("revert_of")
1209 .and_then(Value::as_str)
1210 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1211 self.rollback(
1212 sub,
1213 revert_of,
1214 actor,
1215 observer,
1216 scopes,
1217 &because,
1218 now_ms,
1219 )?;
1220 p = LoopPersisted::from_value(sub.load_state()?)?;
1223 }
1224 }
1225
1226 let applied = AppliedRecord {
1227 applied_at_ms: now_ms,
1228 target_ref: rec.target_ref.clone(),
1229 rollbackable: rec.rollbackable,
1230 created_hashes: created,
1231 metric: rec.metric.clone(),
1232 };
1233 let prev = p.audit_heads.get(rec_hash).cloned();
1234 let audit = AuditRecord {
1235 rec_hash: rec_hash.into(),
1236 from: Some(status),
1237 to: RecStatus::Applied,
1238 actor: actor.into(),
1239 observer_type: observer,
1240 because,
1241 previous_audit_hash: prev,
1242 gating: gating.cloned(),
1243 at_ms: now_ms,
1244 };
1245 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1246 p.audit_heads.insert(rec_hash.into(), audit_hash);
1247 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1248 p.applied.insert(rec_hash.into(), applied.clone());
1249 sub.store_state(&p.to_value()?)?;
1250 Ok(applied)
1251 }
1252
1253 #[allow(clippy::too_many_arguments)]
1256 pub fn rollback<S: OmsSubstrate>(
1257 &self,
1258 sub: &mut S,
1259 rec_hash: &str,
1260 actor: &str,
1261 observer: ObserverType,
1262 scopes: &ScopeSet,
1263 because: &str,
1264 now_ms: i64,
1265 ) -> Result<()> {
1266 if !scopes.has(Scope::Apply) {
1267 return Err(Error::ScopeDenied("apply".into()));
1268 }
1269 let because = validate_because(because)?;
1270 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1271 let status = *p
1272 .status_index
1273 .get(rec_hash)
1274 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1275 if !status.can_transition_to(RecStatus::RolledBack, false) {
1276 return Err(Error::LifecycleViolation(format!(
1277 "{} -> rolled_back",
1278 status.as_str()
1279 )));
1280 }
1281 let applied = p
1282 .applied
1283 .get(rec_hash)
1284 .cloned()
1285 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1286 if !applied.rollbackable {
1287 return Err(Error::LifecycleViolation(
1288 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1289 ));
1290 }
1291 for h in &applied.created_hashes {
1292 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1293 }
1294 let prev = p.audit_heads.get(rec_hash).cloned();
1295 let audit = AuditRecord {
1296 rec_hash: rec_hash.into(),
1297 from: Some(status),
1298 to: RecStatus::RolledBack,
1299 actor: actor.into(),
1300 observer_type: observer,
1301 because,
1302 previous_audit_hash: prev,
1303 gating: None,
1304 at_ms: now_ms,
1305 };
1306 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1307 p.audit_heads.insert(rec_hash.into(), audit_hash);
1308 p.status_index
1309 .insert(rec_hash.into(), RecStatus::RolledBack);
1310 sub.store_state(&p.to_value()?)?;
1311 Ok(())
1312 }
1313
1314 pub fn recommendations<S: OmsSubstrate>(
1319 &self,
1320 sub: &S,
1321 status_filter: Option<RecStatus>,
1322 ) -> Result<Vec<Recommendation>> {
1323 let p = LoopPersisted::from_value(sub.load_state()?)?;
1324 let grains = sub.grains_of_type(
1325 crate::model::grain_type::RECOMMENDATION,
1326 Some(LOOP_NS),
1327 ReadOpts {
1328 live_only: false,
1329 since_ms: None,
1330 },
1331 )?;
1332 let mut out = Vec::new();
1333 for g in grains {
1334 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1335 rec.status = p
1336 .status_index
1337 .get(&g.hash)
1338 .copied()
1339 .unwrap_or(RecStatus::Pending);
1340 if let Some(f) = status_filter {
1341 if rec.status != f {
1342 continue;
1343 }
1344 }
1345 out.push(rec);
1346 }
1347 out.sort_by(|a, b| {
1356 b.severity
1357 .cmp(&a.severity)
1358 .then(a.created_at_ms.cmp(&b.created_at_ms))
1359 .then(a.dedup_key.cmp(&b.dedup_key))
1360 .then(a.hash.cmp(&b.hash))
1361 });
1362 Ok(out)
1363 }
1364
1365 pub fn analyzer_settings<S: OmsSubstrate>(
1368 &self,
1369 sub: &S,
1370 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1371 let p = LoopPersisted::from_value(sub.load_state()?)?;
1372 Ok(self
1373 .analyzers
1374 .iter()
1375 .map(|a| {
1376 let m = a.manifest();
1377 let cfg = p.config.get(&m.id);
1378 crate::config::AnalyzerSetting {
1379 id: m.id.clone(),
1380 title: m.title.clone(),
1381 description: m.description.clone(),
1382 tier: format!("{:?}", m.tier),
1383 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1384 default_on: m.default_on,
1385 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1386 severity_floor: cfg
1387 .and_then(|c| c.severity_floor)
1388 .map(|s| s.as_str().to_string()),
1389 }
1390 })
1391 .collect())
1392 }
1393
1394 pub fn set_analyzer_config<S: OmsSubstrate>(
1400 &self,
1401 sub: &mut S,
1402 analyzer_id: &str,
1403 update: crate::config::AnalyzerConfigUpdate,
1404 scopes: &ScopeSet,
1405 ) -> Result<crate::config::AnalyzerConfig> {
1406 if !scopes.has(Scope::Admin) {
1407 return Err(Error::ScopeDenied("admin".into()));
1408 }
1409 let manifest = self
1410 .analyzers
1411 .iter()
1412 .map(|a| a.manifest())
1413 .find(|m| m.id == analyzer_id)
1414 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1415 if let Some(params) = &update.params {
1417 manifest.resolve_params(params)?;
1418 }
1419 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1420 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1421 if let Some(enabled) = update.enabled {
1422 cfg.enabled = Some(enabled);
1423 }
1424 if update.clear_floor {
1425 cfg.severity_floor = None;
1426 } else if let Some(floor) = update.severity_floor {
1427 cfg.severity_floor = Some(floor);
1428 }
1429 if let Some(params) = update.params {
1430 cfg.params = params;
1431 }
1432 if let Some(ns) = update.namespaces {
1433 cfg.namespaces = ns;
1434 }
1435 let stored = cfg.clone();
1436 sub.store_state(&p.to_value()?)?;
1437 Ok(stored)
1438 }
1439
1440 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1443 let p = LoopPersisted::from_value(sub.load_state()?)?;
1444 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1445 out.sort_by(|a, b| {
1449 a.measured_at_ms
1450 .cmp(&b.measured_at_ms)
1451 .then(a.horizon_ms.cmp(&b.horizon_ms))
1452 .then(a.metric.cmp(&b.metric))
1453 .then(a.rec_hash.cmp(&b.rec_hash))
1454 });
1455 Ok(out)
1456 }
1457
1458 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1462 let p = LoopPersisted::from_value(sub.load_state()?)?;
1463 let (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
1464 let recs = self.recommendations(sub, None)?;
1465 let mut pending = 0;
1466 let mut applied = 0;
1467 for r in &recs {
1468 match r.status {
1469 RecStatus::Pending => pending += 1,
1470 RecStatus::Applied => applied += 1,
1471 _ => {}
1472 }
1473 }
1474 let stale = match p.state.last_run_ms {
1476 None => true,
1477 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1478 };
1479 Ok(Health {
1480 last_run_ms: p.state.last_run_ms,
1481 grains_since_run,
1482 error_events_since_run,
1483 pending,
1484 applied,
1485 total: recs.len() as u64,
1486 stale,
1487 })
1488 }
1489
1490 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1495 let recs = self.recommendations(sub, None)?;
1496 let mut m = LlmMetrics::default();
1497 for r in &recs {
1498 if !matches!(r.origin, Origin::Llm { .. }) {
1499 continue;
1500 }
1501 m.proposed += 1;
1502 match r.status {
1503 RecStatus::Pending => m.pending += 1,
1504 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1505 RecStatus::Rejected => m.rejected += 1,
1506 RecStatus::Expired => {}
1507 }
1508 }
1509 let decided = m.approved + m.rejected;
1510 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1511 Ok(m)
1512 }
1513}
1514
1515#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1517pub struct Health {
1518 #[serde(skip_serializing_if = "Option::is_none")]
1519 pub last_run_ms: Option<i64>,
1520 pub grains_since_run: u64,
1521 pub error_events_since_run: u64,
1522 pub pending: u64,
1523 pub applied: u64,
1524 pub total: u64,
1525 pub stale: bool,
1528}
1529
1530#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1532pub struct LlmMetrics {
1533 pub proposed: u64,
1536 pub pending: u64,
1537 pub approved: u64,
1539 pub rejected: u64,
1540 #[serde(skip_serializing_if = "Option::is_none")]
1542 pub approval_rate: Option<f64>,
1543}
1544
1545fn measure_outcomes<S: OmsSubstrate>(
1552 sub: &S,
1553 p: &mut LoopPersisted,
1554 now_ms: i64,
1555) -> Result<Vec<OutcomeInput>> {
1556 let mut due: Vec<(String, crate::config::AppliedRecord, i64)> = Vec::new();
1558 for (h, a) in &p.applied {
1559 if p.status_index.get(h) != Some(&RecStatus::Applied) {
1560 continue;
1561 }
1562 let Some(metric) = &a.metric else { continue };
1563 let done = p.measured.get(h).cloned().unwrap_or_default();
1564 for horizon in metric.horizons() {
1565 if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
1566 due.push((h.clone(), a.clone(), horizon));
1567 }
1568 }
1569 }
1570
1571 let mut out = Vec::new();
1572 for (rec_hash, applied, horizon) in due {
1573 let metric = applied.metric.as_ref().unwrap();
1574 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1575 continue; };
1577 let regressed = current > metric.baseline + f64::EPSILON;
1578 p.outcomes.entry(rec_hash.clone()).or_default().push(
1579 crate::recommendation::OutcomeResult {
1580 rec_hash: rec_hash.clone(),
1581 metric: metric.metric.clone(),
1582 baseline: metric.baseline,
1583 current,
1584 verdict: if regressed { "regressed" } else { "held" }.into(),
1585 horizon_ms: horizon,
1586 measured_at_ms: now_ms,
1587 },
1588 );
1589 p.measured.entry(rec_hash.clone()).or_default().push(horizon);
1590 if regressed {
1591 out.push(OutcomeInput {
1592 rec_hash,
1593 target_ref: applied.target_ref.clone(),
1594 metric: metric.metric.clone(),
1595 baseline: metric.baseline,
1596 current,
1597 unit: metric.unit.clone(),
1598 });
1599 }
1600 }
1601 Ok(out)
1602}
1603
1604fn measure_metric<S: SubstrateRead>(
1606 sub: &S,
1607 metric: &crate::recommendation::MetricSnapshot,
1608 since_ms: i64,
1609) -> Result<Option<f64>> {
1610 match metric.metric.as_str() {
1611 "tool_error_recurrence" => {
1616 let Some(tool) = &metric.subject else { return Ok(None) };
1617 let tools = sub.grains_of_type(
1618 crate::model::grain_type::TOOL,
1619 None,
1620 ReadOpts { live_only: true, since_ms: Some(since_ms) },
1621 )?;
1622 let n = tools
1623 .iter()
1624 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
1625 .filter(|t| {
1626 metric.relation.as_deref().is_none_or(|sig| {
1629 crate::analyzers::tool_failure::normalize_signature(
1630 t.tool_content().unwrap_or(""),
1631 ) == sig
1632 })
1633 })
1634 .count();
1635 Ok(Some(n as f64))
1636 }
1637 "contradiction_recurrence" => {
1641 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
1642 return Ok(None);
1643 };
1644 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
1645 let distinct: BTreeSet<String> = facts
1646 .iter()
1647 .filter(|f| {
1648 f.fact_relation()
1649 .is_some_and(|r| normalize_ident(r) == *relation)
1650 })
1651 .filter_map(|f| f.fact_object().map(normalize_ident))
1652 .collect();
1653 Ok(Some(distinct.len().saturating_sub(1) as f64))
1654 }
1655 _ => Ok(None),
1656 }
1657}
1658
1659fn scoped_live_facts<S: SubstrateRead>(
1662 sub: &S,
1663 namespace: Option<&str>,
1664 subject: &str,
1665) -> Result<Vec<GrainRecord>> {
1666 let facts = sub.grains_of_type(
1667 crate::model::grain_type::FACT,
1668 None,
1669 ReadOpts { live_only: true, since_ms: None },
1670 )?;
1671 Ok(facts
1672 .into_iter()
1673 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
1674 .filter(|f| {
1675 f.fact_subject()
1676 .is_some_and(|s| normalize_ident(s) == subject)
1677 })
1678 .collect())
1679}
1680
1681fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
1693 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
1694 return false;
1695 };
1696 if fields.is_empty() {
1697 return false;
1698 }
1699 let Ok(Some(grain)) = sub.grain(&target) else {
1700 return false;
1701 };
1702 if grain.valid_to_ms.is_some() {
1712 return false;
1713 }
1714 fields.iter().all(|(k, v)| {
1716 if k == "namespace" {
1717 return v
1718 .as_str()
1719 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
1720 }
1721 match (v, grain.fields.get(k)) {
1722 (Value::String(a), Some(Value::String(b))) => {
1723 normalize_ident(a) == normalize_ident(b)
1724 }
1725 (a, Some(b)) => a == b,
1726 (_, None) => false,
1727 }
1728 })
1729}
1730
1731const MIN_LLM_CONFIDENCE: f64 = 0.75;
1734
1735const DISCOVER_INSTRUCTIONS: &str = "You review an agent's memory for quality. \
1740Given deterministic findings and the evidence they cite, propose ADDITIONAL \
1741findings the deterministic checks would miss (e.g. a semantic contradiction, a \
1742stale assumption, a duplicated meaning). SCORING: propose a finding ONLY if you \
1743are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
1744useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
1745earns 0. When in doubt, propose nothing — an empty list is the correct answer \
1746when there is nothing worth flagging. The 'approved' and 'rejected' lists, when \
1747present, show findings this reviewer recently accepted or rejected — prefer the \
1748kind they accept and avoid the kind they reject. Every proposal MUST cite one \
1749or more evidence hashes from the bundle, target a memory entity, and include \
1750your confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
1751\"target\":\"entity:<ns>/<subject>\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
1752\"confidence\":0.0}]}. Propose nothing you cannot ground in the evidence.";
1753
1754const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
1760fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
1761the facts it relies on are actually present in the cited evidence, NOT that its \
1762conclusion is stated verbatim. Decompose the finding into the factual claims it \
1763depends on. Mark supported=true when those facts are present in the evidence \
1764(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
1765on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
1766different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
1767
1768const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
1770each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
1771never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
1772SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
1773'possible' findings with no concrete defect, and reject any claimed \
1774inconsistency or contradiction that is not backed by at least two actually \
1775conflicting facts in the cited evidence. (2) Context — does the finding \
1776correctly read its cited evidence, or misinterpret what the grains say? Keep a \
1777finding when it names a genuine, specific problem grounded in its evidence and \
1778materially useful to a human reviewer; otherwise reject it, and default to \
1779keep=false when uncertain. Do NOT reject a finding for being 'already known', \
1780redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
1781grounded cross-fact inconsistency with two conflicting facts is exactly what to \
1782KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
1783{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
1784
1785const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
1787guidance note to help a human reviewer decide. Do not restate the finding. Return \
1788JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
1789
1790fn push_evidence(
1793 evidence: &mut Vec<crate::llm::EvidenceItem>,
1794 bundle: &mut BTreeSet<String>,
1795 g: &GrainRecord,
1796) {
1797 if evidence.len() < 64 && bundle.insert(g.hash.clone()) {
1798 evidence.push(crate::llm::EvidenceItem {
1799 hash: g.hash.clone(),
1800 grain_type: g.grain_type.clone(),
1801 text: crate::llm::cap(&grain_brief(g), 400),
1802 });
1803 }
1804}
1805
1806fn grain_brief(g: &GrainRecord) -> String {
1808 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
1809 return format!("{s} {r} {o}");
1810 }
1811 for key in ["content", "body", "text", "summary"] {
1812 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
1813 if !v.is_empty() {
1814 return v.to_string();
1815 }
1816 }
1817 }
1818 String::new()
1819}
1820
1821fn stamp_llm(
1826 model: &str,
1827 d: &crate::llm::LlmDraft,
1828 target_ref: String,
1829 cited: Vec<String>,
1830 confidence: f64,
1831 now_ms: i64,
1832) -> Recommendation {
1833 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
1834 let mut args = serde_json::Map::new();
1835 args.insert("text".into(), Value::from(summary_text));
1836 let guidance = if d.guidance.trim().is_empty() {
1837 None
1838 } else {
1839 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
1840 };
1841 let action = ActionKind::Flag;
1842 let mut data = serde_json::Map::new();
1843 data.insert("source".into(), Value::from("llm"));
1844 Recommendation {
1845 hash: String::new(),
1846 analyzer: "loop.llm/1".to_string(),
1847 params_snapshot: serde_json::Map::new(),
1848 origin: Origin::Llm { model: model.to_string() },
1849 target_ref: target_ref.clone(),
1850 action_kind: action,
1851 dedup_key: dedup_key("llm", &target_ref, action),
1852 summary: Summary::new("llm.discover", args),
1853 severity: Severity::Low,
1854 proposal: Proposal::Data { data },
1855 destructive: false,
1856 rollbackable: false,
1857 evidence: cited,
1858 evidence_query: None,
1859 metric: None,
1860 confidence: confidence.clamp(0.0, 1.0),
1862 importance: 0.3,
1863 created_at_ms: now_ms,
1864 guidance,
1865 evalset_hash: None,
1866 status: RecStatus::Pending,
1867 }
1868}
1869
1870fn stamp(
1871 m: &AnalyzerManifest,
1872 params: &crate::manifest::Params,
1873 d: crate::recommendation::RecDraft,
1874 now_ms: i64,
1875) -> Result<Recommendation> {
1876 let target = TargetRef::parse(&d.target_ref)?;
1877 crate::recommendation::validate_code_rules(
1881 d.action_kind,
1882 target.target_class(),
1883 d.evalset_hash.as_deref(),
1884 )?;
1885 let dedup = dedup_key(m.family(), &d.target_ref, d.action_kind);
1886 let destructive = match &d.proposal {
1887 Proposal::Cal { cal } => cal::contains_destructive(cal),
1888 _ => false,
1889 };
1890 let rollbackable = match &d.proposal {
1891 Proposal::Cal { .. } => !destructive,
1892 Proposal::Edit { .. } => false,
1893 Proposal::Data { .. } => d.action_kind == ActionKind::CodeRevision,
1896 };
1897 let mut evidence = d.evidence;
1898 evidence.truncate(MAX_EVIDENCE);
1899 let origin = match m.trust_class {
1904 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
1905 _ => Origin::Builtin,
1906 };
1907 Ok(Recommendation {
1908 hash: String::new(),
1909 analyzer: m.id.clone(),
1910 params_snapshot: params.snapshot(),
1911 origin,
1912 target_ref: target.as_string(),
1913 action_kind: d.action_kind,
1914 dedup_key: dedup,
1915 summary: d.summary,
1916 severity: d.severity,
1917 proposal: d.proposal,
1918 destructive,
1919 rollbackable,
1920 evidence,
1921 evidence_query: d.evidence_query,
1922 metric: d.metric,
1923 confidence: d.confidence,
1924 importance: d.importance,
1925 created_at_ms: now_ms,
1926 guidance: None,
1927 evalset_hash: d.evalset_hash,
1928 status: RecStatus::Pending,
1929 })
1930}
1931
1932fn validate_because(because: &str) -> Result<String> {
1933 let trimmed = because.trim();
1934 if trimmed.is_empty() {
1935 return Err(Error::InvalidProposal(
1936 "a BECAUSE reason is required".into(),
1937 ));
1938 }
1939 if trimmed.chars().count() > MAX_BECAUSE {
1940 return Err(Error::InvalidProposal(format!(
1941 "BECAUSE exceeds {MAX_BECAUSE} chars"
1942 )));
1943 }
1944 Ok(trimmed.to_string())
1945}
1946
1947fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
1948 for req in &m.requires {
1949 match req {
1950 Capability::Forks if !caps.forks => return Some("forks"),
1951 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
1952 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
1953 _ => {}
1954 }
1955 }
1956 None
1957}
1958
1959fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
1960 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
1961}
1962
1963fn gate(
1964 opts: &RunOptions,
1965 p: &LoopPersisted,
1966 new_grains: u64,
1967 new_errors: u64,
1968 now_ms: i64,
1969) -> Option<SkipReason> {
1970 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
1971 if !any {
1972 return None;
1973 }
1974 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
1975 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
1976 let stale_ok = opts
1977 .if_stale_ms
1978 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
1979 if min_new_ok || min_err_ok || stale_ok {
1980 return None;
1981 }
1982 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
1984 Some(SkipReason::NotStale)
1985 } else {
1986 Some(SkipReason::MinNewNotMet)
1987 }
1988}
1989
1990fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
1991 let opts = ReadOpts {
1992 live_only: false,
1993 since_ms: watermark.map(|w| w + 1),
1994 };
1995 let mut new_grains = 0u64;
1996 let mut new_errors = 0u64;
1997 for t in [
1998 crate::model::grain_type::FACT,
1999 crate::model::grain_type::EVENT,
2000 crate::model::grain_type::TOOL,
2001 crate::model::grain_type::OBSERVATION,
2002 ] {
2003 let g = sub.grains_of_type(t, None, opts)?;
2004 new_grains += g.len() as u64;
2005 if t == crate::model::grain_type::TOOL {
2007 new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
2008 }
2009 }
2010 Ok((new_grains, new_errors))
2011}
2012
2013fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
2014 let grains = sub.grains_of_type(
2015 crate::model::grain_type::RECOMMENDATION,
2016 Some(LOOP_NS),
2017 ReadOpts {
2018 live_only: false,
2019 since_ms: None,
2020 },
2021 )?;
2022 let mut set = BTreeSet::new();
2023 for g in grains {
2024 let status = p
2025 .status_index
2026 .get(&g.hash)
2027 .copied()
2028 .unwrap_or(RecStatus::Pending);
2029 if matches!(
2034 status,
2035 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
2036 ) {
2037 if let Some(key) = g.str_field("dedup_key") {
2038 set.insert(key.to_string());
2039 }
2040 }
2041 }
2042 Ok(set)
2043}
2044
2045fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
2046 let g = sub
2047 .grain(rec_hash)?
2048 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
2049 Recommendation::from_fields(rec_hash, &g.fields)
2050}
2051
2052const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
2055 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
2056 approve it to acknowledge it and let it expire.";
2057
2058const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
2061 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
2062 acknowledge it and let it expire.";
2063
2064pub(crate) fn ensure_executable(proposal: &Proposal) -> Result<()> {
2077 match proposal {
2078 Proposal::Cal { .. } => Ok(()),
2079 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
2080 Proposal::Data { data } => {
2083 if data.get("revert_of").and_then(Value::as_str).is_some() {
2084 Ok(())
2085 } else {
2086 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
2087 }
2088 }
2089 }
2090}