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 #[serde(default, skip_serializing_if = "Option::is_none")]
133 pub llm_funnel: Option<LlmFunnel>,
134}
135
136#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
147pub struct LlmFunnel {
148 pub evidence: u64,
150 pub proposed: u64,
152 pub cited: u64,
154 pub dropped_uncited: u64,
158 pub dropped_target: u64,
161 pub grounded: u64,
163 pub kept: u64,
165 pub stored: u64,
167}
168
169impl RunResult {
170 fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
171 RunResult {
172 outcome: RunOutcome::Skipped,
173 skip_reason: Some(reason),
174 new_grains,
175 new_error_events,
176 proposed: 0,
177 deduped: 0,
178 stored: 0,
179 auto_applied: 0,
180 llm_funnel: None,
181 analyzers_run: vec![],
182 analyzers_skipped: vec![],
183 }
184 }
185
186 pub fn ran(&self) -> bool {
187 self.outcome == RunOutcome::Ran
188 }
189}
190
191pub struct Engine {
194 analyzers: Vec<Box<dyn Analyzer>>,
195 policy: crate::policy::Policy,
196 llm: Option<Box<dyn crate::llm::LlmBackend>>,
199 ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
204}
205
206struct AnalysisPass {
207 survivors: Vec<Recommendation>,
208 proposed: u64,
209 deduped: u64,
210 analyzers_run: Vec<String>,
211 analyzers_skipped: Vec<AnalyzerSkip>,
212 llm_funnel: Option<LlmFunnel>,
213}
214
215impl Engine {
216 pub fn with_builtins() -> Self {
219 Engine {
220 analyzers: crate::analyzer::builtin_analyzers(),
221 policy: crate::policy::Policy::default(),
222 llm: None,
223 ground_llm: None,
224 }
225 }
226
227 pub fn empty() -> Self {
229 Engine {
230 analyzers: vec![],
231 policy: crate::policy::Policy::default(),
232 llm: None,
233 ground_llm: None,
234 }
235 }
236
237 pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
239 self.policy = policy;
240 self
241 }
242
243 pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
248 self.llm = Some(backend);
249 self
250 }
251
252 pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
256 self.ground_llm = Some(backend);
257 self
258 }
259
260 pub fn policy(&self) -> &crate::policy::Policy {
261 &self.policy
262 }
263
264 pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
266 self.analyzers.push(analyzer);
267 }
268
269 pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
270 &self.analyzers
271 }
272
273 pub fn analyze_only<S: OmsSubstrate>(
282 &self,
283 sub: &S,
284 opts: &RunOptions,
285 overrides: &BTreeMap<String, Map<String, Value>>,
286 now_ms: i64,
287 ) -> Result<Vec<Recommendation>> {
288 let persisted = LoopPersisted::from_value(sub.load_state()?)?;
289 let analysis_watermark = if opts.full_sweep {
290 None
291 } else {
292 persisted.state.watermark_ms
293 };
294 Ok(self
295 .analysis_pass(
296 sub,
297 &persisted,
298 opts,
299 overrides,
300 analysis_watermark,
301 now_ms,
302 &[],
303 )?
304 .survivors)
305 }
306
307 pub fn run<S: OmsSubstrate>(
310 &self,
311 sub: &mut S,
312 opts: &RunOptions,
313 now_ms: i64,
314 ) -> Result<RunResult> {
315 let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
316 let watermark = persisted.state.watermark_ms;
317 let analysis_watermark = if opts.full_sweep { None } else { watermark };
322
323 let (new_grains, new_error_events) = count_new(sub, watermark)?;
324 if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
325 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
326 }
327
328 let outcome_inputs = measure_outcomes(sub, &mut persisted, now_ms)?;
331
332 let AnalysisPass {
333 survivors,
334 proposed,
335 deduped,
336 analyzers_run,
337 analyzers_skipped,
338 llm_funnel,
339 } = self.analysis_pass(
340 &*sub,
341 &persisted,
342 opts,
343 &BTreeMap::new(),
344 analysis_watermark,
345 now_ms,
346 &outcome_inputs,
347 )?;
348
349 let mut stored = 0u64;
352 let mut auto_applied = 0u64;
353 for mut rec in survivors {
354 let spec = rec.to_grain_spec(LOOP_NS)?;
355 let hash = sub.put_grain(&spec)?;
356 rec.hash = hash.clone();
357 let actor = format!("engine:{}", rec.analyzer);
358 let audit = AuditRecord {
359 rec_hash: hash.clone(),
360 from: None,
361 to: RecStatus::Pending,
362 actor: actor.clone(),
363 observer_type: ObserverType::System,
364 because: "analyzer proposed".into(),
365 previous_audit_hash: None,
366 gating: None,
367 at_ms: now_ms,
368 };
369 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
370 persisted
371 .status_index
372 .insert(hash.clone(), RecStatus::Pending);
373 persisted.creators.insert(hash.clone(), actor);
374 if !matches!(rec.origin, Origin::Builtin) {
379 if let Some(trigger) = &opts.triggering_actor {
380 persisted.co_creators.insert(hash.clone(), trigger.clone());
381 }
382 }
383 persisted.audit_heads.insert(hash.clone(), audit_hash);
384 stored += 1;
385
386 if self.can_auto_apply(&*sub, &rec) {
387 self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
388 auto_applied += 1;
389 }
390 }
391
392 persisted.state.last_run_ms = Some(now_ms);
393 persisted.state.watermark_ms = Some(now_ms);
394 sub.store_state(&persisted.to_value()?)?;
395
396 Ok(RunResult {
397 outcome: RunOutcome::Ran,
398 skip_reason: None,
399 new_grains,
400 new_error_events,
401 proposed,
402 deduped,
403 stored,
404 auto_applied,
405 analyzers_run,
406 analyzers_skipped,
407 llm_funnel,
408 })
409 }
410
411 #[allow(clippy::too_many_arguments)]
415 fn analysis_pass<S: OmsSubstrate>(
416 &self,
417 sub: &S,
418 persisted: &LoopPersisted,
419 opts: &RunOptions,
420 external_overrides: &BTreeMap<String, Map<String, Value>>,
421 analysis_watermark: Option<i64>,
422 now_ms: i64,
423 outcome_inputs: &[OutcomeInput],
424 ) -> Result<AnalysisPass> {
425 let existing = existing_dedup_keys(sub, persisted)?;
426 let mut analyzers_run = Vec::new();
427 let mut analyzers_skipped = Vec::new();
428 let mut candidates: Vec<Recommendation> = Vec::new();
429 let caps = sub.capabilities();
430
431 for analyzer in &self.analyzers {
432 let m = analyzer.manifest();
433 let cfg = persisted.config.get(&m.id);
434 let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
435 if !enabled {
436 analyzers_skipped.push(AnalyzerSkip {
437 id: m.id.clone(),
438 reason: "disabled".into(),
439 });
440 continue;
441 }
442 if self.policy.denies(m.family()) {
443 analyzers_skipped.push(AnalyzerSkip {
444 id: m.id.clone(),
445 reason: "denied by host policy".into(),
446 });
447 continue;
448 }
449 if let Some(missing) = missing_capability(m, caps) {
450 analyzers_skipped.push(AnalyzerSkip {
451 id: m.id.clone(),
452 reason: format!("missing capability: {missing}"),
453 });
454 continue;
455 }
456 let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
457 if let Some(extra) = external_overrides.get(&m.id) {
458 for (key, value) in extra {
459 param_overrides.insert(key.clone(), value.clone());
460 }
461 }
462 let params = match m.resolve_params(¶m_overrides) {
463 Ok(p) => p,
464 Err(e) => {
465 analyzers_skipped.push(AnalyzerSkip {
466 id: m.id.clone(),
467 reason: e.to_string(),
468 });
469 continue;
470 }
471 };
472 let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
473 let ns_slice: &[String] = if ns_owned.is_empty() {
474 &opts.namespaces
475 } else {
476 &ns_owned
477 };
478 let reader: &dyn SubstrateRead = sub;
479 let ctx = AnalyzeCtx::new(
480 reader,
481 ¶ms,
482 ns_slice,
483 analysis_watermark,
484 now_ms,
485 outcome_inputs,
486 );
487 match analyzer.analyze(&ctx) {
488 Ok(drafts) => {
489 analyzers_run.push(m.id.clone());
490 for draft in drafts {
491 match stamp(m, ¶ms, draft, now_ms) {
492 Ok(rec) => candidates.push(rec),
493 Err(e) => analyzers_skipped.push(AnalyzerSkip {
494 id: m.id.clone(),
495 reason: e.to_string(),
496 }),
497 }
498 }
499 }
500 Err(e) => analyzers_skipped.push(AnalyzerSkip {
501 id: m.id.clone(),
502 reason: e.to_string(),
503 }),
504 }
505 }
506
507 let mut funnel = LlmFunnel::default();
508 if self.llm.is_some() {
509 candidates.extend(self.discover(
510 sub,
511 &candidates,
512 analysis_watermark,
513 &opts.namespaces,
514 now_ms,
515 &mut funnel,
516 ));
517 }
518
519 let proposed = candidates.len() as u64;
520 let mut seen = BTreeSet::new();
521 let mut survivors = Vec::new();
522 for candidate in candidates {
523 let family = crate::manifest::analyzer_family(&candidate.analyzer);
524 let floor = [
525 severity_floor_for(persisted, &candidate.analyzer),
526 self.policy.severity_floor(family),
527 ]
528 .into_iter()
529 .flatten()
530 .max();
531 if floor.is_some_and(|floor| candidate.severity < floor) {
532 continue;
533 }
534 if !seen.insert(candidate.dedup_key.clone()) {
535 continue;
536 }
537 if existing.contains(&candidate.dedup_key) {
538 continue;
539 }
540 if persisted
541 .cooldowns
542 .get(&candidate.dedup_key)
543 .is_some_and(|until| now_ms < *until)
544 {
545 continue;
546 }
547 survivors.push(candidate);
548 }
549 let deduped = proposed - survivors.len() as u64;
550 if self.llm.is_some() {
551 self.enrich(&mut survivors);
552 }
553 Ok(AnalysisPass {
554 survivors,
555 proposed,
556 deduped,
557 analyzers_run,
558 analyzers_skipped,
559 llm_funnel: self.llm.is_some().then_some(funnel),
560 })
561 }
562
563 fn discover<S: OmsSubstrate>(
570 &self,
571 sub: &S,
572 candidates: &[Recommendation],
573 watermark: Option<i64>,
574 namespaces: &[String],
575 now_ms: i64,
576 funnel: &mut LlmFunnel,
577 ) -> Vec<Recommendation> {
578 let Some(llm) = &self.llm else {
579 return Vec::new();
580 };
581 let findings: Vec<crate::llm::FindingBrief> = candidates
582 .iter()
583 .take(32)
584 .map(|c| crate::llm::FindingBrief {
585 analyzer: c.analyzer.clone(),
586 summary: c.summary.render(),
587 target: c.target_ref.clone(),
588 severity: c.severity.as_str().to_string(),
589 })
590 .collect();
591 let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
597 let mut bundle: BTreeSet<String> = BTreeSet::new();
598 let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
599 'cited: for c in candidates {
600 for h in &c.evidence {
601 if evidence.len() >= CITED_SEED_CAP {
610 break 'cited;
611 }
612 if !bundle.contains(h) {
613 if let Ok(Some(g)) = sub.grain(h) {
614 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
615 }
616 }
617 }
618 }
619 let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
620 vec![None]
621 } else {
622 namespaces.iter().map(|n| Some(n.as_str())).collect()
623 };
624 let opts = ReadOpts { live_only: true, since_ms: watermark };
625 let mut tool_seeded = 0usize;
637 'tools: for ns in &scan_ns {
638 if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
639 for g in recent {
640 if tool_seeded >= TOOL_SEED_CAP
641 || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
642 {
643 break 'tools;
644 }
645 if !g.is_error() {
646 continue;
647 }
648 let before = evidence.len();
649 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
650 if evidence.len() > before {
651 tool_seeded += 1;
652 }
653 }
654 }
655 }
656 'notes: for ns in &scan_ns {
666 if let Ok(recent) =
667 sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
668 {
669 for g in recent {
670 if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
671 break 'notes;
672 }
673 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
674 }
675 }
676 }
677 'seed: for gt in [
678 crate::model::grain_type::FACT,
679 crate::model::grain_type::OBSERVATION,
680 ] {
681 for ns in &scan_ns {
682 if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
683 for g in recent {
684 if evidence.len() >= EVIDENCE_CAP {
685 break 'seed;
686 }
687 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
688 }
689 }
690 }
691 }
692 funnel.evidence = evidence.len() as u64;
693 if evidence.is_empty() {
694 return Vec::new(); }
696 let (approved, rejected) = self.llm_history(sub);
701 let request = crate::llm::LlmRequest {
702 loop_proto: 1,
703 op: "discover",
704 instructions: DISCOVER_INSTRUCTIONS,
705 findings: findings.clone(),
706 evidence: evidence.clone(),
707 rejected,
708 approved,
709 };
710 let Ok(body) = serde_json::to_string(&request) else {
711 return Vec::new();
712 };
713 let raw = match llm.complete(&body) {
714 Ok(r) => r,
715 Err(_) => return Vec::new(), };
717 let caps = sub.capabilities();
721 let mut validated: Vec<ValidatedDraft> = Vec::new();
722 let drafts: Vec<_> = crate::llm::parse_discover(&raw)
723 .recommendations
724 .into_iter()
725 .take(crate::llm::MAX_LLM_DRAFTS)
726 .collect();
727 funnel.proposed = drafts.len() as u64;
728 for d in drafts {
729 let cited: Vec<String> =
730 d.evidence.iter().filter(|h| bundle.contains(*h)).cloned().collect();
731 if cited.is_empty() {
732 funnel.dropped_uncited += 1;
733 continue; }
735 let Ok(target) = TargetRef::parse(&d.target) else {
736 funnel.dropped_target += 1;
737 continue;
738 };
739 let tc = target.target_class();
740 if !matches!(tc, "memory" | "query" | "code") {
746 funnel.dropped_target += 1;
747 continue;
748 }
749 let resolved = resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps);
752 if tc == "code" && resolved.is_none() {
757 funnel.dropped_target += 1;
758 continue;
759 }
760 validated.push(ValidatedDraft {
761 draft: d,
762 target_ref: target.as_string(),
763 cited,
764 resolved,
765 });
766 }
767 funnel.cited = validated.len() as u64;
768 if validated.is_empty() {
769 return Vec::new();
770 }
771 let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
778 self.verify_drafts(&**llm, ground, validated, &evidence, now_ms, funnel)
779 }
780
781 fn verify_drafts(
788 &self,
789 llm: &dyn crate::llm::LlmBackend,
790 ground: &dyn crate::llm::LlmBackend,
791 validated: Vec<ValidatedDraft>,
792 evidence: &[crate::llm::EvidenceItem],
793 now_ms: i64,
794 funnel: &mut LlmFunnel,
795 ) -> Vec<Recommendation> {
796 use crate::llm::*;
797 let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
798 evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
799 let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
800 cited
801 .iter()
802 .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
803 .collect()
804 };
805
806 let claims: Vec<GroundItem> = validated
810 .iter()
811 .enumerate()
812 .map(|(i, v)| GroundItem {
813 id: i,
814 claim: claim_text(&v.draft, v.resolved.as_ref()),
815 evidence: ev_for(&v.cited),
816 })
817 .collect();
818 let ground_req = GroundRequest {
819 loop_proto: 1,
820 op: "ground",
821 instructions: GROUND_INSTRUCTIONS,
822 claims,
823 };
824 let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
825 .ok()
826 .and_then(|b| ground.complete(&b).ok())
827 {
828 Some(raw) => parse_ground(&raw)
829 .results
830 .into_iter()
831 .filter(|r| r.supported)
832 .map(|r| r.id)
833 .collect(),
834 None => return Vec::new(),
835 };
836 funnel.grounded = grounded.len() as u64;
837 if grounded.is_empty() {
838 return Vec::new();
839 }
840
841 let items: Vec<VerifyItem> = validated
847 .iter()
848 .enumerate()
849 .filter(|(i, _)| grounded.contains(i))
850 .map(|(i, v)| VerifyItem {
851 id: i,
852 summary: claim_text(&v.draft, v.resolved.as_ref()),
854 target: v.target_ref.clone(),
855 evidence: ev_for(&v.cited),
856 })
857 .collect();
858 let verify_req = VerifyRequest {
859 loop_proto: 1,
860 op: "verify",
861 instructions: VERIFY_INSTRUCTIONS,
862 findings: items,
863 };
864 let verdicts: std::collections::BTreeMap<usize, f64> =
865 match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
866 Some(raw) => parse_verify(&raw)
867 .results
868 .into_iter()
869 .filter(|r| r.keep)
870 .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
871 .collect(),
872 None => return Vec::new(),
873 };
874
875 funnel.kept = verdicts.len() as u64;
876 let mut out = Vec::new();
880 for (i, v) in validated.into_iter().enumerate() {
881 if let Some(&conf) = verdicts.get(&i) {
882 if conf >= MIN_LLM_CONFIDENCE {
883 out.push(stamp_llm(
884 llm.model(),
885 &v.draft,
886 v.target_ref,
887 v.cited,
888 v.resolved,
889 conf,
890 now_ms,
891 ));
892 }
893 }
894 }
895 funnel.stored = out.len() as u64;
896 out
897 }
898
899 fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
904 const MAX: usize = 20;
905 let Ok(mut recs) = self.recommendations(sub, None) else {
906 return (Vec::new(), Vec::new());
907 };
908 recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
909 recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
910 let mut approved = Vec::new();
911 let mut rejected = Vec::new();
912 for r in &recs {
913 match r.status {
914 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
915 if approved.len() < MAX =>
916 {
917 approved.push(r.summary.render());
918 }
919 RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
920 _ => {}
921 }
922 }
923 (approved, rejected)
924 }
925
926 fn enrich(&self, survivors: &mut [Recommendation]) {
931 let Some(llm) = &self.llm else {
932 return;
933 };
934 if survivors.is_empty() {
935 return;
936 }
937 let findings: Vec<crate::llm::FindingBrief> = survivors
938 .iter()
939 .map(|r| crate::llm::FindingBrief {
940 analyzer: r.analyzer.clone(),
941 summary: r.summary.render(),
942 target: r.target_ref.clone(),
943 severity: r.severity.as_str().to_string(),
944 })
945 .collect();
946 let request = crate::llm::LlmRequest {
947 loop_proto: 1,
948 op: "enrich",
949 instructions: ENRICH_INSTRUCTIONS,
950 findings,
951 evidence: Vec::new(),
952 rejected: Vec::new(),
953 approved: Vec::new(),
954 };
955 let Ok(body) = serde_json::to_string(&request) else {
956 return;
957 };
958 let raw = match llm.complete(&body) {
959 Ok(r) => r,
960 Err(_) => return,
961 };
962 for note in crate::llm::parse_enrich(&raw).notes {
963 if note.guidance.trim().is_empty() {
964 continue;
965 }
966 if let Some(r) = survivors
967 .iter_mut()
968 .find(|r| r.target_ref == note.target && r.guidance.is_none())
969 {
970 r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
971 }
972 }
973 }
974
975 fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
984 if !rec.origin.auto_apply_eligible() || rec.destructive {
985 return false;
986 }
987 let manifest_ok = self
991 .analyzers
992 .iter()
993 .map(|a| a.manifest())
994 .find(|m| m.id == rec.analyzer)
995 .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
996 if !manifest_ok {
997 return false;
998 }
999 let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1000 return false;
1001 };
1002 let family = crate::manifest::analyzer_family(&rec.analyzer);
1003 if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1004 return false;
1005 }
1006 match &rec.proposal {
1011 Proposal::Cal { cal } => cal
1012 .lines()
1013 .map(str::trim)
1014 .filter(|l| !l.is_empty())
1015 .all(|l| supersede_is_value_identical(sub, l)),
1016 _ => false,
1017 }
1018 }
1019
1020 fn auto_apply<S: OmsSubstrate>(
1023 &self,
1024 sub: &mut S,
1025 p: &mut LoopPersisted,
1026 rec: &Recommendation,
1027 now_ms: i64,
1028 ) -> Result<()> {
1029 let mut created = Vec::new();
1030 if let Proposal::Cal { cal } = &rec.proposal {
1031 if cal.lines().map(str::trim).any(is_definition_statement) {
1036 return Err(Error::InvalidProposal(
1037 "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1038 auto-applied: it changes what every future context contains, so it \
1039 requires a human APPROVE + APPLY with BECAUSE"
1040 .into(),
1041 ));
1042 }
1043 for r in sub.execute_cal(cal)? {
1044 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1045 created.push(h.to_string());
1046 }
1047 }
1048 }
1049 let applied = AppliedRecord {
1050 applied_at_ms: now_ms,
1051 target_ref: rec.target_ref.clone(),
1052 rollbackable: rec.rollbackable,
1053 created_hashes: created,
1054 inverse_cal: None,
1055 metric: rec.metric.clone(),
1056 };
1057 let prev = p.audit_heads.get(&rec.hash).cloned();
1058 let audit = AuditRecord {
1059 rec_hash: rec.hash.clone(),
1060 from: Some(RecStatus::Pending),
1061 to: RecStatus::Applied,
1062 actor: "policy:auto".into(),
1063 observer_type: ObserverType::Policy,
1064 because: "auto-applied per host policy".into(),
1065 previous_audit_hash: prev,
1066 gating: None,
1067 at_ms: now_ms,
1068 };
1069 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1070 p.audit_heads.insert(rec.hash.clone(), audit_hash);
1071 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1072 p.applied.insert(rec.hash.clone(), applied);
1073 Ok(())
1074 }
1075
1076 #[allow(clippy::too_many_arguments)]
1079 pub fn review<S: OmsSubstrate>(
1080 &self,
1081 sub: &mut S,
1082 rec_hash: &str,
1083 decision: Decision,
1084 actor: &str,
1085 observer: ObserverType,
1086 scopes: &ScopeSet,
1087 because: &str,
1088 now_ms: i64,
1089 ) -> Result<()> {
1090 if !scopes.has(Scope::Review) {
1091 return Err(Error::ScopeDenied("review".into()));
1092 }
1093 let because = validate_because(because)?;
1094 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1095 let status = *p
1096 .status_index
1097 .get(rec_hash)
1098 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1099 let to = match decision {
1100 Decision::Approve => RecStatus::Approved,
1101 Decision::Reject => RecStatus::Rejected,
1102 };
1103 if !status.can_transition_to(to, false) {
1104 return Err(Error::LifecycleViolation(format!(
1105 "{} -> {}",
1106 status.as_str(),
1107 to.as_str()
1108 )));
1109 }
1110 if to == RecStatus::Approved {
1111 if let Some(creator) = p.creators.get(rec_hash) {
1112 if creator == actor {
1113 return Err(Error::SelfApproval(format!(
1114 "{actor} created this recommendation"
1115 )));
1116 }
1117 }
1118 if let Some(trigger) = p.co_creators.get(rec_hash) {
1119 if trigger == actor {
1120 return Err(Error::SelfApproval(format!(
1121 "{actor} triggered the run that authored this recommendation"
1122 )));
1123 }
1124 }
1125 }
1126 let prev = p.audit_heads.get(rec_hash).cloned();
1127 let audit = AuditRecord {
1128 rec_hash: rec_hash.into(),
1129 from: Some(status),
1130 to,
1131 actor: actor.into(),
1132 observer_type: observer,
1133 because,
1134 previous_audit_hash: prev,
1135 gating: None,
1136 at_ms: now_ms,
1137 };
1138 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1139 p.audit_heads.insert(rec_hash.into(), audit_hash);
1140 p.status_index.insert(rec_hash.into(), to);
1141 if to == RecStatus::Rejected {
1142 if let Ok(rec) = load_rec(sub, rec_hash) {
1143 const BASE_MS: i64 = 7 * 86_400_000;
1148 const CAP_MS: i64 = 90 * 86_400_000;
1149 let strikes = p.cooldown_strikes.entry(rec.dedup_key.clone()).or_insert(0);
1150 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
1151 *strikes = strikes.saturating_add(1);
1152 p.cooldowns.insert(rec.dedup_key, now_ms + interval);
1153 }
1154 }
1155 sub.store_state(&p.to_value()?)?;
1156 Ok(())
1157 }
1158
1159 pub fn preflight_apply<S: OmsSubstrate>(
1176 &self,
1177 sub: &S,
1178 rec_hash: &str,
1179 scopes: &ScopeSet,
1180 allow_destructive: bool,
1181 has_gating: bool,
1182 ) -> Result<()> {
1183 if !scopes.has(Scope::Apply) {
1184 return Err(Error::ScopeDenied("apply".into()));
1185 }
1186 let rec = load_rec(sub, rec_hash)?;
1187 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1188 return Err(Error::DestructiveGated(
1189 "destructive apply requires admin scope + allow_destructive".into(),
1190 ));
1191 }
1192 ensure_executable(rec.action_kind, &rec.proposal)?;
1193 if requires_gating(rec.action_kind) && !has_gating {
1194 return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1195 }
1196 Ok(())
1197 }
1198
1199 #[allow(clippy::too_many_arguments)]
1203 pub fn apply<S: OmsSubstrate>(
1204 &self,
1205 sub: &mut S,
1206 rec_hash: &str,
1207 actor: &str,
1208 observer: ObserverType,
1209 scopes: &ScopeSet,
1210 because: &str,
1211 allow_destructive: bool,
1212 now_ms: i64,
1213 ) -> Result<AppliedRecord> {
1214 self.apply_inner(
1215 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1216 )
1217 }
1218
1219 pub fn gating_evidence<S: OmsSubstrate>(
1226 &self,
1227 sub: &S,
1228 rec_hash: &str,
1229 run_id: &str,
1230 ) -> Result<crate::recommendation::GatingEvidence> {
1231 let rec = self
1232 .recommendations(sub, None)?
1233 .into_iter()
1234 .find(|r| r.hash == rec_hash)
1235 .ok_or_else(|| {
1236 Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1237 })?;
1238 let pin = rec.evalset_hash.ok_or_else(|| {
1239 Error::InvalidProposal(
1240 "this recommendation pins no evalset — a gating run applies only \
1241 to code and adapter revisions"
1242 .into(),
1243 )
1244 })?;
1245 match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1250 Some(run) => Ok(crate::recommendation::GatingEvidence {
1251 evalset_hash: pin,
1252 run_id: run.run_id,
1253 passed: run.passed,
1254 failed: run.failed,
1255 }),
1256 None => Err(Error::InvalidProposal(format!(
1257 "no recorded gate run '{run_id}' for evalset {pin} — run \
1258 `areev eval run --evalset {pin} ...` first"
1259 ))),
1260 }
1261 }
1262
1263 #[allow(clippy::too_many_arguments)]
1267 pub fn apply_gated<S: OmsSubstrate>(
1268 &self,
1269 sub: &mut S,
1270 rec_hash: &str,
1271 actor: &str,
1272 observer: ObserverType,
1273 scopes: &ScopeSet,
1274 because: &str,
1275 allow_destructive: bool,
1276 gating: &crate::recommendation::GatingEvidence,
1277 now_ms: i64,
1278 ) -> Result<AppliedRecord> {
1279 self.apply_inner(
1280 sub,
1281 rec_hash,
1282 actor,
1283 observer,
1284 scopes,
1285 because,
1286 allow_destructive,
1287 Some(gating),
1288 now_ms,
1289 )
1290 }
1291
1292 #[allow(clippy::too_many_arguments)]
1293 fn apply_inner<S: OmsSubstrate>(
1294 &self,
1295 sub: &mut S,
1296 rec_hash: &str,
1297 actor: &str,
1298 observer: ObserverType,
1299 scopes: &ScopeSet,
1300 because: &str,
1301 allow_destructive: bool,
1302 gating: Option<&crate::recommendation::GatingEvidence>,
1303 now_ms: i64,
1304 ) -> Result<AppliedRecord> {
1305 if !scopes.has(Scope::Apply) {
1306 return Err(Error::ScopeDenied("apply".into()));
1307 }
1308 let because = validate_because(because)?;
1309 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1310 let status = *p
1311 .status_index
1312 .get(rec_hash)
1313 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1314 if !status.can_transition_to(RecStatus::Applied, false) {
1315 return Err(Error::LifecycleViolation(format!(
1316 "{} -> applied (approve first)",
1317 status.as_str()
1318 )));
1319 }
1320 let rec = load_rec(sub, rec_hash)?;
1321 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1322 return Err(Error::DestructiveGated(
1323 "destructive apply requires admin scope + allow_destructive".into(),
1324 ));
1325 }
1326 if requires_gating(rec.action_kind) {
1331 let g = gating
1332 .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1333 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1334 if g.evalset_hash != pin {
1335 return Err(Error::InvalidProposal(format!(
1336 "gating ran evalset {} but the recommendation is pinned \
1337 to {pin} (Rule E1)",
1338 g.evalset_hash
1339 )));
1340 }
1341 match sub.grain(pin)? {
1342 Some(evalset) if evalset.is_live() => {}
1343 Some(_) => {
1344 return Err(Error::InvalidProposal(
1345 "the pinned evalset was superseded after gating — \
1346 the recommendation must re-gate (Rule E1)"
1347 .into(),
1348 ))
1349 }
1350 None => {
1351 return Err(Error::InvalidProposal(format!(
1352 "pinned evalset {pin} not found in the substrate"
1353 )))
1354 }
1355 }
1356 if g.failed > 0 {
1357 return Err(Error::InvalidProposal(format!(
1358 "the gating run failed {}/{} cases — a failing gate \
1359 admits nothing",
1360 g.failed,
1361 g.passed + g.failed
1362 )));
1363 }
1364 }
1365
1366 let mut created = Vec::new();
1368 let mut inverse_cal: Option<String> = None;
1372 match &rec.proposal {
1373 Proposal::Cal { cal } => {
1374 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1375 if !is_definition_statement(line) {
1376 continue;
1377 }
1378 match sub.definition_inverse(line)? {
1379 Some(inv) => inverse_cal = Some(inv),
1380 None => {
1381 return Err(Error::InvalidProposal(format!(
1382 "this substrate cannot record a rollback inverse for {line:?}; \
1383 a definition rewrite that ROLLBACK could not undo is refused \
1384 rather than applied"
1385 )))
1386 }
1387 }
1388 }
1389 let rows = sub.execute_cal(cal)?;
1390 for r in rows {
1391 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1392 created.push(h.to_string());
1393 }
1394 }
1395 }
1396 Proposal::Edit { .. } => {
1399 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1400 }
1401 Proposal::Data { data } if requires_gating(rec.action_kind) => {
1409 let relation = if rec.action_kind == ActionKind::AdapterRevision {
1410 "mg:adapter_promotion"
1411 } else {
1412 "mg:code_promotion"
1413 };
1414 let mut promoted = data.clone();
1422 if let Some(Value::String(src)) = promoted.remove("source") {
1423 let address = sub.put_blob(src.as_bytes())?;
1424 promoted.insert("code_address".into(), Value::from(address));
1425 }
1426 let mut spec = crate::substrate::GrainSpec::new(
1427 crate::model::grain_type::FACT,
1428 LOOP_NS,
1429 )
1430 .with_field("subject", rec.target_ref.clone())
1431 .with_field("relation", relation)
1432 .with_field(
1433 "object",
1434 serde_json::to_string(&promoted)
1435 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1436 )
1437 .with_field("rec_hash", rec_hash.to_string());
1438 if let Some(g) = gating {
1439 spec = spec
1440 .with_field("gating_evalset", g.evalset_hash.clone())
1441 .with_field("gating_run_id", g.run_id.clone());
1442 }
1443 created.push(sub.put_grain(&spec)?);
1444 }
1445 Proposal::Data { data } => {
1446 let revert_of = data
1451 .get("revert_of")
1452 .and_then(Value::as_str)
1453 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1454 self.rollback(
1455 sub,
1456 revert_of,
1457 actor,
1458 observer,
1459 scopes,
1460 &because,
1461 now_ms,
1462 )?;
1463 p = LoopPersisted::from_value(sub.load_state()?)?;
1466 }
1467 }
1468
1469 let applied = AppliedRecord {
1470 applied_at_ms: now_ms,
1471 target_ref: rec.target_ref.clone(),
1472 rollbackable: rec.rollbackable,
1473 created_hashes: created,
1474 inverse_cal,
1475 metric: rec.metric.clone(),
1476 };
1477 let prev = p.audit_heads.get(rec_hash).cloned();
1478 let audit = AuditRecord {
1479 rec_hash: rec_hash.into(),
1480 from: Some(status),
1481 to: RecStatus::Applied,
1482 actor: actor.into(),
1483 observer_type: observer,
1484 because,
1485 previous_audit_hash: prev,
1486 gating: gating.cloned(),
1487 at_ms: now_ms,
1488 };
1489 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1490 p.audit_heads.insert(rec_hash.into(), audit_hash);
1491 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1492 p.applied.insert(rec_hash.into(), applied.clone());
1493 sub.store_state(&p.to_value()?)?;
1494 Ok(applied)
1495 }
1496
1497 #[allow(clippy::too_many_arguments)]
1500 pub fn rollback<S: OmsSubstrate>(
1501 &self,
1502 sub: &mut S,
1503 rec_hash: &str,
1504 actor: &str,
1505 observer: ObserverType,
1506 scopes: &ScopeSet,
1507 because: &str,
1508 now_ms: i64,
1509 ) -> Result<()> {
1510 if !scopes.has(Scope::Apply) {
1511 return Err(Error::ScopeDenied("apply".into()));
1512 }
1513 let because = validate_because(because)?;
1514 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1515 let status = *p
1516 .status_index
1517 .get(rec_hash)
1518 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1519 if !status.can_transition_to(RecStatus::RolledBack, false) {
1520 return Err(Error::LifecycleViolation(format!(
1521 "{} -> rolled_back",
1522 status.as_str()
1523 )));
1524 }
1525 let applied = p
1526 .applied
1527 .get(rec_hash)
1528 .cloned()
1529 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1530 if !applied.rollbackable {
1531 return Err(Error::LifecycleViolation(
1532 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1533 ));
1534 }
1535 for h in &applied.created_hashes {
1536 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1537 }
1538 if let Some(inverse) = &applied.inverse_cal {
1545 sub.execute_cal(inverse)?;
1546 }
1547 let prev = p.audit_heads.get(rec_hash).cloned();
1548 let audit = AuditRecord {
1549 rec_hash: rec_hash.into(),
1550 from: Some(status),
1551 to: RecStatus::RolledBack,
1552 actor: actor.into(),
1553 observer_type: observer,
1554 because,
1555 previous_audit_hash: prev,
1556 gating: None,
1557 at_ms: now_ms,
1558 };
1559 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1560 p.audit_heads.insert(rec_hash.into(), audit_hash);
1561 p.status_index
1562 .insert(rec_hash.into(), RecStatus::RolledBack);
1563 sub.store_state(&p.to_value()?)?;
1564 Ok(())
1565 }
1566
1567 pub fn recommendations<S: OmsSubstrate>(
1572 &self,
1573 sub: &S,
1574 status_filter: Option<RecStatus>,
1575 ) -> Result<Vec<Recommendation>> {
1576 let p = LoopPersisted::from_value(sub.load_state()?)?;
1577 let grains = sub.grains_of_type(
1578 crate::model::grain_type::RECOMMENDATION,
1579 Some(LOOP_NS),
1580 ReadOpts {
1581 live_only: false,
1582 since_ms: None,
1583 },
1584 )?;
1585 let mut out = Vec::new();
1586 for g in grains {
1587 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1588 rec.status = p
1589 .status_index
1590 .get(&g.hash)
1591 .copied()
1592 .unwrap_or(RecStatus::Pending);
1593 if let Some(f) = status_filter {
1594 if rec.status != f {
1595 continue;
1596 }
1597 }
1598 out.push(rec);
1599 }
1600 out.sort_by(|a, b| {
1609 b.severity
1610 .cmp(&a.severity)
1611 .then(a.created_at_ms.cmp(&b.created_at_ms))
1612 .then(a.dedup_key.cmp(&b.dedup_key))
1613 .then(a.hash.cmp(&b.hash))
1614 });
1615 Ok(out)
1616 }
1617
1618 pub fn analyzer_settings<S: OmsSubstrate>(
1621 &self,
1622 sub: &S,
1623 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1624 let p = LoopPersisted::from_value(sub.load_state()?)?;
1625 Ok(self
1626 .analyzers
1627 .iter()
1628 .map(|a| {
1629 let m = a.manifest();
1630 let cfg = p.config.get(&m.id);
1631 crate::config::AnalyzerSetting {
1632 id: m.id.clone(),
1633 title: m.title.clone(),
1634 description: m.description.clone(),
1635 tier: format!("{:?}", m.tier),
1636 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1637 default_on: m.default_on,
1638 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1639 severity_floor: cfg
1640 .and_then(|c| c.severity_floor)
1641 .map(|s| s.as_str().to_string()),
1642 }
1643 })
1644 .collect())
1645 }
1646
1647 pub fn set_analyzer_config<S: OmsSubstrate>(
1653 &self,
1654 sub: &mut S,
1655 analyzer_id: &str,
1656 update: crate::config::AnalyzerConfigUpdate,
1657 scopes: &ScopeSet,
1658 ) -> Result<crate::config::AnalyzerConfig> {
1659 if !scopes.has(Scope::Admin) {
1660 return Err(Error::ScopeDenied("admin".into()));
1661 }
1662 let manifest = self
1663 .analyzers
1664 .iter()
1665 .map(|a| a.manifest())
1666 .find(|m| m.id == analyzer_id)
1667 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1668 if let Some(params) = &update.params {
1670 manifest.resolve_params(params)?;
1671 }
1672 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1673 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1674 if let Some(enabled) = update.enabled {
1675 cfg.enabled = Some(enabled);
1676 }
1677 if update.clear_floor {
1678 cfg.severity_floor = None;
1679 } else if let Some(floor) = update.severity_floor {
1680 cfg.severity_floor = Some(floor);
1681 }
1682 if let Some(params) = update.params {
1683 cfg.params = params;
1684 }
1685 if let Some(ns) = update.namespaces {
1686 cfg.namespaces = ns;
1687 }
1688 let stored = cfg.clone();
1689 sub.store_state(&p.to_value()?)?;
1690 Ok(stored)
1691 }
1692
1693 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1696 let p = LoopPersisted::from_value(sub.load_state()?)?;
1697 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1698 out.sort_by(|a, b| {
1702 a.measured_at_ms
1703 .cmp(&b.measured_at_ms)
1704 .then(a.horizon_ms.cmp(&b.horizon_ms))
1705 .then(a.metric.cmp(&b.metric))
1706 .then(a.rec_hash.cmp(&b.rec_hash))
1707 });
1708 Ok(out)
1709 }
1710
1711 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1715 let p = LoopPersisted::from_value(sub.load_state()?)?;
1716 let (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
1717 let recs = self.recommendations(sub, None)?;
1718 let mut pending = 0;
1719 let mut applied = 0;
1720 for r in &recs {
1721 match r.status {
1722 RecStatus::Pending => pending += 1,
1723 RecStatus::Applied => applied += 1,
1724 _ => {}
1725 }
1726 }
1727 let stale = match p.state.last_run_ms {
1729 None => true,
1730 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1731 };
1732 Ok(Health {
1733 last_run_ms: p.state.last_run_ms,
1734 grains_since_run,
1735 error_events_since_run,
1736 pending,
1737 applied,
1738 total: recs.len() as u64,
1739 stale,
1740 })
1741 }
1742
1743 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1748 let recs = self.recommendations(sub, None)?;
1749 let mut m = LlmMetrics::default();
1750 for r in &recs {
1751 if !matches!(r.origin, Origin::Llm { .. }) {
1752 continue;
1753 }
1754 m.proposed += 1;
1755 match r.status {
1756 RecStatus::Pending => m.pending += 1,
1757 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1758 RecStatus::Rejected => m.rejected += 1,
1759 RecStatus::Expired => {}
1760 }
1761 }
1762 let decided = m.approved + m.rejected;
1763 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1764 Ok(m)
1765 }
1766}
1767
1768#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1770pub struct Health {
1771 #[serde(skip_serializing_if = "Option::is_none")]
1772 pub last_run_ms: Option<i64>,
1773 pub grains_since_run: u64,
1774 pub error_events_since_run: u64,
1775 pub pending: u64,
1776 pub applied: u64,
1777 pub total: u64,
1778 pub stale: bool,
1781}
1782
1783#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1785pub struct LlmMetrics {
1786 pub proposed: u64,
1789 pub pending: u64,
1790 pub approved: u64,
1792 pub rejected: u64,
1793 #[serde(skip_serializing_if = "Option::is_none")]
1795 pub approval_rate: Option<f64>,
1796}
1797
1798fn measure_outcomes<S: OmsSubstrate>(
1805 sub: &S,
1806 p: &mut LoopPersisted,
1807 now_ms: i64,
1808) -> Result<Vec<OutcomeInput>> {
1809 let mut due: Vec<(String, crate::config::AppliedRecord, i64)> = Vec::new();
1811 for (h, a) in &p.applied {
1812 if p.status_index.get(h) != Some(&RecStatus::Applied) {
1813 continue;
1814 }
1815 let Some(metric) = &a.metric else { continue };
1816 let done = p.measured.get(h).cloned().unwrap_or_default();
1817 for horizon in metric.horizons() {
1818 if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
1819 due.push((h.clone(), a.clone(), horizon));
1820 }
1821 }
1822 }
1823
1824 let mut out = Vec::new();
1825 for (rec_hash, applied, horizon) in due {
1826 let metric = applied.metric.as_ref().unwrap();
1827 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1828 continue; };
1830 let regressed = crate::recommendation::is_regression(
1831 metric.baseline,
1832 current,
1833 metric.higher_is_better,
1834 );
1835 p.outcomes.entry(rec_hash.clone()).or_default().push(
1836 crate::recommendation::OutcomeResult {
1837 rec_hash: rec_hash.clone(),
1838 metric: metric.metric.clone(),
1839 baseline: metric.baseline,
1840 current,
1841 verdict: if regressed { "regressed" } else { "held" }.into(),
1842 horizon_ms: horizon,
1843 measured_at_ms: now_ms,
1844 },
1845 );
1846 p.measured.entry(rec_hash.clone()).or_default().push(horizon);
1847 if regressed {
1848 out.push(OutcomeInput {
1849 rec_hash,
1850 target_ref: applied.target_ref.clone(),
1851 metric: metric.metric.clone(),
1852 baseline: metric.baseline,
1853 current,
1854 unit: metric.unit.clone(),
1855 higher_is_better: metric.higher_is_better,
1856 });
1857 }
1858 }
1859 Ok(out)
1860}
1861
1862pub(crate) fn measure_metric<S: SubstrateRead>(
1864 sub: &S,
1865 metric: &crate::recommendation::MetricSnapshot,
1866 since_ms: i64,
1867) -> Result<Option<f64>> {
1868 match metric.metric.as_str() {
1869 "tool_error_recurrence" => {
1874 let Some(tool) = &metric.subject else { return Ok(None) };
1875 let tools = sub.grains_of_type(
1876 crate::model::grain_type::TOOL,
1877 None,
1878 ReadOpts { live_only: true, since_ms: Some(since_ms) },
1879 )?;
1880 let n = tools
1881 .iter()
1882 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
1883 .filter(|t| {
1884 metric.relation.as_deref().is_none_or(|sig| {
1887 crate::analyzers::tool_failure::normalize_signature(
1888 t.tool_content().unwrap_or(""),
1889 ) == sig
1890 })
1891 })
1892 .count();
1893 Ok(Some(n as f64))
1894 }
1895 "contradiction_recurrence" => {
1899 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
1900 return Ok(None);
1901 };
1902 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
1903 let distinct: BTreeSet<String> = facts
1904 .iter()
1905 .filter(|f| {
1906 f.fact_relation()
1907 .is_some_and(|r| normalize_ident(r) == *relation)
1908 })
1909 .filter_map(|f| f.fact_object().map(normalize_ident))
1910 .collect();
1911 Ok(Some(distinct.len().saturating_sub(1) as f64))
1912 }
1913 m if m.starts_with("evalset:") => {
1926 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
1927 return Ok(None);
1928 };
1929 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
1930 return Ok(None);
1931 };
1932 Ok(match field {
1936 "failed" => Some(run.failed as f64),
1937 "passed" => Some(run.passed as f64),
1938 "total" => Some(run.total() as f64),
1939 "error_rate" => match run.total() {
1940 0 => None, t => Some(run.failed as f64 / t as f64),
1942 },
1943 other => run.field(other),
1944 })
1945 }
1946 _ => Ok(None),
1947 }
1948}
1949
1950fn scoped_live_facts<S: SubstrateRead>(
1953 sub: &S,
1954 namespace: Option<&str>,
1955 subject: &str,
1956) -> Result<Vec<GrainRecord>> {
1957 let facts = sub.grains_of_type(
1958 crate::model::grain_type::FACT,
1959 None,
1960 ReadOpts { live_only: true, since_ms: None },
1961 )?;
1962 Ok(facts
1963 .into_iter()
1964 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
1965 .filter(|f| {
1966 f.fact_subject()
1967 .is_some_and(|s| normalize_ident(s) == subject)
1968 })
1969 .collect())
1970}
1971
1972fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
1984 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
1985 return false;
1986 };
1987 if fields.is_empty() {
1988 return false;
1989 }
1990 let Ok(Some(grain)) = sub.grain(&target) else {
1991 return false;
1992 };
1993 if grain.valid_to_ms.is_some() {
2003 return false;
2004 }
2005 fields.iter().all(|(k, v)| {
2007 if k == "namespace" {
2008 return v
2009 .as_str()
2010 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2011 }
2012 match (v, grain.fields.get(k)) {
2013 (Value::String(a), Some(Value::String(b))) => {
2014 normalize_ident(a) == normalize_ident(b)
2015 }
2016 (a, Some(b)) => a == b,
2017 (_, None) => false,
2018 }
2019 })
2020}
2021
2022const EVIDENCE_CAP: usize = 64;
2032const CITED_SEED_CAP: usize = 24;
2033const TOOL_SEED_CAP: usize = 16;
2034const NOTE_SEED_CAP: usize = 8;
2039const LENS_RESERVE: usize = 24;
2040
2041const MIN_LLM_CONFIDENCE: f64 = 0.75;
2044
2045const DISCOVER_INSTRUCTIONS: &str = "You review an agent's memory for quality. \
2050Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2051findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2052stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2053recurring cost or hand-off the agent's own setup could remove). \
2054The deterministic findings already cover what the ERROR TEXT says; restating \
2055one of them earns nothing. The evidence may also contain OUTCOME records — a \
2056run's observable shape together with whether it was accepted or rejected. A \
2057problem that raised no error at all is exactly the kind the deterministic \
2058checks cannot see, so compare the rejected outcomes against the accepted \
2059ones: a feature they share and the accepted ones lack is a candidate rule. \
2060Require at least two rejected outcomes before proposing one — a single \
2061rejection is an anecdote, not a pattern. \
2062SCORING: propose a finding ONLY if you \
2063are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2064useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2065earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2066when there is nothing worth flagging. The 'approved' and 'rejected' lists, when \
2067present, show findings this reviewer recently accepted or rejected — prefer the \
2068kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2069or more evidence hashes from the bundle, name a 'target', and include your \
2070confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
2071\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
2072\"confidence\":0.0,\"proposal\":{...}}]}. \
2073OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2074you are not asking to change anything. Include it ONLY when the evidence \
2075supports a specific change, choosing exactly one kind: \
2076(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2077\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) \
2078preventing a recurring mistake, phrased to apply BEFORE it happens (e.g. \
2079'Refund a subscription before cancelling it; refunds on cancelled \
2080subscriptions are refused'). \
2081(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2082entity target — a durable fact the agent keeps having to be told (an alias, a \
2083settled default, a preference). 'relation' is a short identifier (letters, \
2084digits, _ - . :), not a sentence. \
2085(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2086\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2087assembles the agent's context, when the evidence shows it retrieves the wrong \
2088things. Give the FULL new body; it replaces the old one. \
2089(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2090with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2091workflow. Only these paths are editable: 'edges.<i>.cond', \
2092'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2093holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2094(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2095— full replacement source for that tool. It is applied only after a recorded \
2096evaluation run passes, so propose one only when the evidence shows the current \
2097code is the defect. \
2098The subject of a fact, the name of a query, the plan hash and the tool name \
2099all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2100becomes a change a human reviewer may apply, so it must be fully supported by \
2101the cited evidence. Propose nothing you cannot ground in the evidence.";
2102
2103const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2109fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2110the facts it relies on are actually present in the cited evidence, NOT that its \
2111conclusion is stated verbatim. Decompose the finding into the factual claims it \
2112depends on. Mark supported=true when those facts are present in the evidence \
2113(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2114on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2115different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2116
2117const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2119each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2120never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2121SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2122'possible' findings with no concrete defect, and reject any claimed \
2123inconsistency or contradiction that is not backed by at least two actually \
2124conflicting facts in the cited evidence. (2) Context — does the finding \
2125correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2126finding when it names a genuine, specific problem grounded in its evidence and \
2127materially useful to a human reviewer; otherwise reject it, and default to \
2128keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2129redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2130grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2131KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2132{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2133
2134const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2136guidance note to help a human reviewer decide. Do not restate the finding. Return \
2137JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2138
2139fn push_evidence(
2145 evidence: &mut Vec<crate::llm::EvidenceItem>,
2146 bundle: &mut BTreeSet<String>,
2147 ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2148 g: &GrainRecord,
2149) {
2150 if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2151 ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2152 evidence.push(crate::llm::EvidenceItem {
2153 hash: g.hash.clone(),
2154 grain_type: g.grain_type.clone(),
2155 text: crate::llm::cap(&grain_brief(g), 400),
2156 });
2157 }
2158}
2159
2160fn grain_brief(g: &GrainRecord) -> String {
2162 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2163 return format!("{s} {r} {o}");
2164 }
2165 if let Some(t) = g.tool_name() {
2170 let status = if g.is_error() { "error" } else { "ok" };
2171 let out = g.tool_content().unwrap_or("");
2172 return format!("tool {t} {status}: {out}");
2173 }
2174 for key in ["content", "body", "text", "summary"] {
2175 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2176 if !v.is_empty() {
2177 return v.to_string();
2178 }
2179 }
2180 }
2181 String::new()
2182}
2183
2184fn sanitize_lesson(s: &str) -> String {
2189 sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2190}
2191
2192fn sanitize_line(s: &str, max: usize) -> String {
2198 let cleaned: String =
2199 s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2200 crate::llm::cap(cleaned.trim(), max)
2201}
2202
2203fn sanitize_relation(s: &str) -> Option<String> {
2207 let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2208 if r.is_empty()
2209 || !r
2210 .chars()
2211 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2212 {
2213 return None;
2214 }
2215 Some(r)
2216}
2217
2218fn safe_definition_body(body: &str) -> bool {
2234 if body.contains('{') || body.contains('}') {
2235 return false;
2236 }
2237 !body
2238 .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2239 .any(|tok| {
2240 ["FORGET", "PURGE", "DROP", "DEFINE"]
2241 .iter()
2242 .any(|kw| tok.eq_ignore_ascii_case(kw))
2243 })
2244}
2245
2246fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2251 let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2252 match resolved {
2253 Some(r) => format!("{summary} {}", r.rendered),
2254 None => summary,
2255 }
2256}
2257
2258struct ValidatedDraft {
2264 draft: crate::llm::LlmDraft,
2265 target_ref: String,
2266 cited: Vec<String>,
2267 resolved: Option<ResolvedProposal>,
2268}
2269
2270struct ResolvedProposal {
2274 action: ActionKind,
2275 proposal: Proposal,
2276 rendered: String,
2279 summary_key: &'static str,
2280 summary_args: serde_json::Map<String, Value>,
2281 rollbackable: bool,
2282 evalset_hash: Option<String>,
2283 importance: f64,
2284 fact_fields: Option<serde_json::Map<String, Value>>,
2288}
2289
2290fn plan_edit_allowed(path: &str) -> bool {
2300 let seg: Vec<&str> = path.split('.').collect();
2301 match seg.as_slice() {
2302 ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2303 ["retries", node] => !node.is_empty(),
2304 _ => false,
2305 }
2306}
2307
2308fn plan_get(body: &Value, path: &str) -> Value {
2311 let mut cur = body;
2312 for seg in path.split('.') {
2313 cur = match cur {
2314 Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2315 Some(v) => v,
2316 None => return Value::Null,
2317 },
2318 Value::Object(o) => match o.get(seg) {
2319 Some(v) => v,
2320 None => return Value::Null,
2321 },
2322 _ => return Value::Null,
2323 };
2324 }
2325 cur.clone()
2326}
2327
2328fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2331 let segs: Vec<&str> = path.split('.').collect();
2332 let Some((last, parents)) = segs.split_last() else {
2333 return false;
2334 };
2335 let mut cur = body;
2336 for seg in parents {
2337 cur = match cur {
2338 Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2339 Some(v) => v,
2340 None => return false,
2341 },
2342 Value::Object(o) => match o.get_mut(*seg) {
2343 Some(v) => v,
2344 None => return false,
2345 },
2346 _ => return false,
2347 };
2348 }
2349 match cur {
2350 Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2351 Some(slot) => {
2352 *slot = to;
2353 true
2354 }
2355 None => false,
2356 },
2357 Value::Object(o) => {
2358 o.insert((*last).to_string(), to);
2359 true
2360 }
2361 _ => false,
2362 }
2363}
2364
2365fn plan_value_ok(path: &str, to: &Value) -> bool {
2370 let seg: Vec<&str> = path.split('.').collect();
2371 match seg.as_slice() {
2372 ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
2373 !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
2374 }),
2375 ["edges", _, "max_cycles"] | ["retries", _] => {
2376 to.as_u64().is_some_and(|n| n <= 1_000)
2377 }
2378 _ => false,
2379 }
2380}
2381
2382fn resolve_proposal<S: OmsSubstrate>(
2388 sub: &S,
2389 d: &crate::llm::LlmDraft,
2390 target: &TargetRef,
2391 cited: &[String],
2392 ns_by_hash: &std::collections::BTreeMap<String, String>,
2393 caps: Capabilities,
2394) -> Option<ResolvedProposal> {
2395 use crate::llm::DraftProposal as P;
2396 let mut args = serde_json::Map::new();
2397 match d.parsed_proposal()? {
2398 P::Lesson { lesson } => {
2400 let lesson = sanitize_lesson(&lesson);
2401 if lesson.is_empty() {
2402 return None;
2403 }
2404 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
2405 args.insert("lesson".into(), Value::from(lesson.clone()));
2406 Some(ResolvedProposal {
2407 action: ActionKind::ClusterFailure,
2411 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2412 rendered: format!("Proposed lesson to record: \"{lesson}\""),
2413 summary_key: "llm.lesson",
2414 summary_args: args,
2415 rollbackable: true,
2416 evalset_hash: None,
2417 importance: 0.5,
2418 fact_fields: Some(fields),
2419 })
2420 }
2421 P::Fact { relation, object } => {
2423 let relation = sanitize_relation(&relation)?;
2424 let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
2425 if object.is_empty() {
2426 return None;
2427 }
2428 let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
2429 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
2430 args.insert("relation".into(), Value::from(relation.clone()));
2431 args.insert("object".into(), Value::from(object.clone()));
2432 Some(ResolvedProposal {
2433 action: ActionKind::Record,
2434 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2435 rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
2436 summary_key: "llm.fact",
2437 summary_args: args,
2438 rollbackable: true,
2439 evalset_hash: None,
2440 importance: 0.5,
2441 fact_fields: Some(fields),
2442 })
2443 }
2444 P::QueryRevision { body } => {
2446 let name = target.opaque();
2447 if name.is_empty()
2451 || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
2452 {
2453 return None;
2454 }
2455 let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
2456 if body.is_empty() || !safe_definition_body(&body) {
2457 return None;
2458 }
2459 let stmt = match target.scheme() {
2460 "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
2461 "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
2462 _ => return None,
2463 };
2464 sub.validate_cal(&stmt).ok()?;
2469 sub.definition_inverse(&stmt).ok().flatten()?;
2470 args.insert("name".into(), Value::from(name));
2471 args.insert("body".into(), Value::from(body.clone()));
2472 Some(ResolvedProposal {
2473 action: ActionKind::Revise,
2474 proposal: Proposal::Cal { cal: stmt },
2475 rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
2476 summary_key: "llm.query_revision",
2477 summary_args: args,
2478 rollbackable: true,
2479 evalset_hash: None,
2480 importance: 0.6,
2481 fact_fields: None,
2482 })
2483 }
2484 P::PlanRevision { edits } => {
2486 if !caps.plans
2487 || target.scheme() != "grain"
2488 || edits.is_empty()
2489 || edits.len() > crate::llm::MAX_PLAN_EDITS
2490 {
2491 return None;
2492 }
2493 let hash = target.opaque();
2494 let g = sub.grain(hash).ok().flatten()?;
2495 if g.grain_type != "workflow" || !g.is_live() {
2496 return None;
2497 }
2498 let mut body = Value::Object(g.fields.clone());
2499 let mut deltas = Vec::new();
2500 let nodes: std::collections::BTreeSet<String> = body
2501 .get("nodes")
2502 .and_then(Value::as_array)
2503 .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
2504 .unwrap_or_default();
2505 for e in &edits {
2506 if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
2507 return None;
2508 }
2509 if let Some(node) = e.path.strip_prefix("retries.") {
2514 if !nodes.contains(node) {
2515 return None;
2516 }
2517 }
2518 if plan_get(&body, &e.path) != e.from {
2521 return None;
2522 }
2523 if e.from == e.to {
2527 return None;
2528 }
2529 if !plan_set(&mut body, &e.path, e.to.clone()) {
2530 return None;
2531 }
2532 deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
2533 }
2534 sub.validate_plan(&body).ok()?;
2538 let Value::Object(fields) = body else {
2539 return None;
2540 };
2541 let stmt = cal::supersede(hash, "workflow", &fields);
2542 sub.validate_cal(&stmt).ok()?;
2547 args.insert("plan".into(), Value::from(hash));
2548 args.insert("edits".into(), Value::from(deltas.join("; ")));
2549 Some(ResolvedProposal {
2550 action: ActionKind::Revise,
2551 proposal: Proposal::Cal { cal: stmt },
2552 rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
2553 summary_key: "llm.plan_revision",
2554 summary_args: args,
2555 rollbackable: true,
2556 evalset_hash: None,
2557 importance: 0.7,
2558 fact_fields: None,
2559 })
2560 }
2561 P::CodeRevision { source } => {
2563 if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
2564 return None;
2565 }
2566 if source.chars().count() > crate::llm::MAX_CODE_LEN {
2567 return None;
2568 }
2569 let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
2573 let mut data = serde_json::Map::new();
2574 data.insert("tool".into(), Value::from(target.opaque()));
2575 data.insert("source".into(), Value::from(source.clone()));
2576 args.insert("tool".into(), Value::from(target.opaque()));
2577 args.insert("bytes".into(), Value::from(source.len() as u64));
2578 Some(ResolvedProposal {
2579 action: ActionKind::CodeRevision,
2580 proposal: Proposal::Data { data },
2581 rendered: format!(
2582 "Proposed new source for tool {} ({} bytes), gated by evalset {}",
2583 target.opaque(),
2584 source.len(),
2585 evalset
2586 ),
2587 summary_key: "llm.code_revision",
2588 summary_args: args,
2589 rollbackable: true,
2590 evalset_hash: Some(evalset),
2591 importance: 0.8,
2592 fact_fields: None,
2593 })
2594 }
2595 }
2596}
2597
2598fn stamp_llm(
2608 model: &str,
2609 d: &crate::llm::LlmDraft,
2610 target_ref: String,
2611 cited: Vec<String>,
2612 resolved: Option<ResolvedProposal>,
2613 confidence: f64,
2614 now_ms: i64,
2615) -> Recommendation {
2616 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2617 let guidance = if d.guidance.trim().is_empty() {
2618 None
2619 } else {
2620 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
2621 };
2622 let (action, proposal, summary, rollbackable, importance, evalset_hash) = match resolved {
2623 Some(mut r) => {
2624 if let Some(mut fields) = r.fact_fields.take() {
2627 fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
2628 r.proposal = Proposal::Cal { cal: cal::add("fact", &fields) };
2629 }
2630 let mut args = r.summary_args;
2631 args.insert("text".into(), Value::from(summary_text));
2632 (
2633 r.action,
2634 r.proposal,
2635 Summary::new(r.summary_key, args),
2636 r.rollbackable,
2637 r.importance,
2638 r.evalset_hash,
2639 )
2640 }
2641 None => {
2642 let mut args = serde_json::Map::new();
2643 args.insert("text".into(), Value::from(summary_text));
2644 let mut data = serde_json::Map::new();
2645 data.insert("source".into(), Value::from("llm"));
2646 (
2647 ActionKind::Flag,
2648 Proposal::Data { data },
2649 Summary::new("llm.discover", args),
2650 false,
2651 0.3,
2652 None,
2653 )
2654 }
2655 };
2656 Recommendation {
2657 hash: String::new(),
2658 analyzer: "loop.llm/1".to_string(),
2659 params_snapshot: serde_json::Map::new(),
2660 origin: Origin::Llm { model: model.to_string() },
2661 target_ref: target_ref.clone(),
2662 action_kind: action,
2663 dedup_key: dedup_key("llm", &target_ref, action),
2664 summary,
2665 severity: Severity::Low,
2666 proposal,
2667 destructive: false,
2668 rollbackable,
2669 evidence: cited,
2670 evidence_query: None,
2671 metric: None,
2672 confidence: confidence.clamp(0.0, 1.0),
2674 importance,
2675 created_at_ms: now_ms,
2676 guidance,
2677 evalset_hash,
2678 status: RecStatus::Pending,
2679 }
2680}
2681
2682fn derived_fact_fields(
2690 target: &TargetRef,
2691 relation: &str,
2692 object: &str,
2693 cited: &[String],
2694 ns_by_hash: &std::collections::BTreeMap<String, String>,
2695) -> Option<serde_json::Map<String, Value>> {
2696 if target.scheme() != "entity" {
2697 return None;
2698 }
2699 let subject = target
2700 .opaque()
2701 .rsplit_once('/')
2702 .map(|(_, s)| s)
2703 .unwrap_or(target.opaque());
2704 if subject.is_empty() {
2705 return None;
2706 }
2707 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
2708 for h in cited {
2709 if let Some(ns) = ns_by_hash.get(h) {
2710 if !ns.is_empty() {
2711 *ns_counts.entry(ns.as_str()).or_default() += 1;
2712 }
2713 }
2714 }
2715 let lesson_ns = ns_counts
2716 .iter()
2717 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
2718 .map(|(ns, _)| ns.to_string());
2719 let mut fields = serde_json::Map::new();
2720 fields.insert("subject".into(), Value::from(subject));
2721 fields.insert("relation".into(), Value::from(relation));
2722 fields.insert("object".into(), Value::from(object));
2723 if let Some(ns) = lesson_ns {
2727 fields.insert("namespace".into(), Value::from(ns));
2728 }
2729 Some(fields)
2730}
2731
2732fn requires_gating(kind: ActionKind) -> bool {
2737 matches!(
2738 kind,
2739 ActionKind::CodeRevision | ActionKind::AdapterRevision
2740 )
2741}
2742
2743fn stamp(
2744 m: &AnalyzerManifest,
2745 params: &crate::manifest::Params,
2746 d: crate::recommendation::RecDraft,
2747 now_ms: i64,
2748) -> Result<Recommendation> {
2749 let target = TargetRef::parse(&d.target_ref)?;
2750 crate::recommendation::validate_code_rules(
2754 d.action_kind,
2755 target.target_class(),
2756 d.evalset_hash.as_deref(),
2757 )?;
2758 let dedup = dedup_key(m.family(), &d.target_ref, d.action_kind);
2759 let destructive = match &d.proposal {
2760 Proposal::Cal { cal } => cal::contains_destructive(cal),
2761 _ => false,
2762 };
2763 let rollbackable = match &d.proposal {
2764 Proposal::Cal { .. } => !destructive,
2765 Proposal::Edit { .. } => false,
2766 Proposal::Data { .. } => requires_gating(d.action_kind),
2770 };
2771 let mut evidence = d.evidence;
2772 evidence.truncate(MAX_EVIDENCE);
2773 let origin = match m.trust_class {
2778 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
2779 _ => Origin::Builtin,
2780 };
2781 Ok(Recommendation {
2782 hash: String::new(),
2783 analyzer: m.id.clone(),
2784 params_snapshot: params.snapshot(),
2785 origin,
2786 target_ref: target.as_string(),
2787 action_kind: d.action_kind,
2788 dedup_key: dedup,
2789 summary: d.summary,
2790 severity: d.severity,
2791 proposal: d.proposal,
2792 destructive,
2793 rollbackable,
2794 evidence,
2795 evidence_query: d.evidence_query,
2796 metric: d.metric,
2797 confidence: d.confidence,
2798 importance: d.importance,
2799 created_at_ms: now_ms,
2800 guidance: None,
2801 evalset_hash: d.evalset_hash,
2802 status: RecStatus::Pending,
2803 })
2804}
2805
2806fn validate_because(because: &str) -> Result<String> {
2807 let trimmed = because.trim();
2808 if trimmed.is_empty() {
2809 return Err(Error::InvalidProposal(
2810 "a BECAUSE reason is required".into(),
2811 ));
2812 }
2813 if trimmed.chars().count() > MAX_BECAUSE {
2814 return Err(Error::InvalidProposal(format!(
2815 "BECAUSE exceeds {MAX_BECAUSE} chars"
2816 )));
2817 }
2818 Ok(trimmed.to_string())
2819}
2820
2821fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
2822 for req in &m.requires {
2823 match req {
2824 Capability::Forks if !caps.forks => return Some("forks"),
2825 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
2826 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
2827 _ => {}
2828 }
2829 }
2830 None
2831}
2832
2833fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
2834 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
2835}
2836
2837fn gate(
2838 opts: &RunOptions,
2839 p: &LoopPersisted,
2840 new_grains: u64,
2841 new_errors: u64,
2842 now_ms: i64,
2843) -> Option<SkipReason> {
2844 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
2845 if !any {
2846 return None;
2847 }
2848 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
2849 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
2850 let stale_ok = opts
2851 .if_stale_ms
2852 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
2853 if min_new_ok || min_err_ok || stale_ok {
2854 return None;
2855 }
2856 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
2858 Some(SkipReason::NotStale)
2859 } else {
2860 Some(SkipReason::MinNewNotMet)
2861 }
2862}
2863
2864fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
2865 let opts = ReadOpts {
2866 live_only: false,
2867 since_ms: watermark.map(|w| w + 1),
2868 };
2869 let mut new_grains = 0u64;
2870 let mut new_errors = 0u64;
2871 for t in [
2872 crate::model::grain_type::FACT,
2873 crate::model::grain_type::EVENT,
2874 crate::model::grain_type::TOOL,
2875 crate::model::grain_type::OBSERVATION,
2876 ] {
2877 let g = sub.grains_of_type(t, None, opts)?;
2878 new_grains += g.len() as u64;
2879 if t == crate::model::grain_type::TOOL {
2881 new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
2882 }
2883 }
2884 Ok((new_grains, new_errors))
2885}
2886
2887fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
2888 let grains = sub.grains_of_type(
2889 crate::model::grain_type::RECOMMENDATION,
2890 Some(LOOP_NS),
2891 ReadOpts {
2892 live_only: false,
2893 since_ms: None,
2894 },
2895 )?;
2896 let mut set = BTreeSet::new();
2897 for g in grains {
2898 let status = p
2899 .status_index
2900 .get(&g.hash)
2901 .copied()
2902 .unwrap_or(RecStatus::Pending);
2903 if matches!(
2908 status,
2909 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
2910 ) {
2911 if let Some(key) = g.str_field("dedup_key") {
2912 set.insert(key.to_string());
2913 }
2914 }
2915 }
2916 Ok(set)
2917}
2918
2919fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
2920 let g = sub
2921 .grain(rec_hash)?
2922 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
2923 Recommendation::from_fields(rec_hash, &g.fields)
2924}
2925
2926pub(crate) fn is_definition_statement(line: &str) -> bool {
2934 let up = line.trim_start().to_ascii_uppercase();
2935 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
2936}
2937
2938const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
2941 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
2942 approve it to acknowledge it and let it expire.";
2943
2944const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
2949 (evalset hash + run id + stats) — use apply_gated";
2950
2951const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
2952 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
2953 acknowledge it and let it expire.";
2954
2955pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
2968 match proposal {
2969 Proposal::Cal { .. } => Ok(()),
2970 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
2971 Proposal::Data { data } => {
2977 if requires_gating(action_kind)
2978 || data.get("revert_of").and_then(Value::as_str).is_some()
2979 {
2980 Ok(())
2981 } else {
2982 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
2983 }
2984 }
2985 }
2986}
2987
2988#[cfg(test)]
2989mod definition_body_tests {
2990 use super::safe_definition_body;
2991
2992 #[test]
2993 fn ordinary_bodies_pass() {
2994 assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
2995 assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
2996 }
2997
2998 #[test]
2999 fn a_body_cannot_close_its_own_block_or_carry_destruction() {
3000 assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
3005 assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
3006 assert!(!safe_definition_body("RECALL facts FORGET abc"));
3008 assert!(!safe_definition_body("recall facts purge older than 1d"));
3009 assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
3010 assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
3012 assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
3014 }
3015}
3016
3017#[cfg(test)]
3018mod plan_edit_tests {
3019 use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
3020 use serde_json::json;
3021
3022 fn plan() -> serde_json::Value {
3023 json!({
3024 "nodes": ["fetch", "review", "post"],
3025 "edges": [
3026 {"src": "fetch", "dst": "review"},
3027 {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
3028 ],
3029 "bindings": {"fetch": "sha256:tool1"},
3030 "retries": {"fetch": 1}
3031 })
3032 }
3033
3034 #[test]
3035 fn the_allowlist_admits_thresholds_and_refuses_topology() {
3036 assert!(plan_edit_allowed("edges.1.cond"));
3037 assert!(plan_edit_allowed("edges.1.max_cycles"));
3038 assert!(plan_edit_allowed("retries.fetch"));
3039 for path in [
3042 "nodes",
3043 "nodes.0",
3044 "edges.0.src",
3045 "edges.0.dst",
3046 "edges",
3047 "bindings.fetch",
3048 "edges.x.cond",
3049 "",
3050 ] {
3051 assert!(!plan_edit_allowed(path), "{path} must not be editable");
3052 }
3053 }
3054
3055 #[test]
3056 fn values_are_type_checked_against_the_field() {
3057 assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
3060 assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
3061 assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
3062 assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
3063 assert!(plan_value_ok("retries.fetch", &json!(3)));
3064 assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
3065 assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
3066 assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
3067 assert!(!plan_value_ok("edges.0.src", &json!("other")));
3068 }
3069
3070 #[test]
3071 fn get_reads_through_arrays_and_objects_and_absence_is_null() {
3072 let p = plan();
3073 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
3074 assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
3075 assert_eq!(plan_get(&p, "retries.review"), json!(null));
3078 assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
3079 assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
3080 }
3081
3082 #[test]
3083 fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
3084 let mut p = plan();
3085 assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
3086 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
3087 assert!(plan_set(&mut p, "retries.review", json!(2)));
3088 assert_eq!(plan_get(&p, "retries.review"), json!(2));
3089 assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
3090 assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
3091 }
3092}
3093
3094#[cfg(test)]
3095mod definition_proposal_tests {
3096 use super::is_definition_statement;
3097
3098 #[test]
3099 fn definition_statements_are_recognized_in_both_spellings() {
3100 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
3101 assert!(is_definition_statement(" define template foo AS { x }"));
3102 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
3103 assert!(!is_definition_statement("ADD fact {}"));
3105 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
3106 assert!(!is_definition_statement("FORGET abc"));
3107 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
3111 }
3112}