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 ground_verdicts: u64,
168 pub ground_call_failed: bool,
172 pub kept: u64,
174 pub stored: u64,
176}
177
178impl RunResult {
179 fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
180 RunResult {
181 outcome: RunOutcome::Skipped,
182 skip_reason: Some(reason),
183 new_grains,
184 new_error_events,
185 proposed: 0,
186 deduped: 0,
187 stored: 0,
188 auto_applied: 0,
189 llm_funnel: None,
190 analyzers_run: vec![],
191 analyzers_skipped: vec![],
192 }
193 }
194
195 pub fn ran(&self) -> bool {
196 self.outcome == RunOutcome::Ran
197 }
198}
199
200pub struct Engine {
203 analyzers: Vec<Box<dyn Analyzer>>,
204 policy: crate::policy::Policy,
205 llm: Option<Box<dyn crate::llm::LlmBackend>>,
208 ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
213}
214
215struct AnalysisPass {
216 survivors: Vec<Recommendation>,
217 proposed: u64,
218 deduped: u64,
219 analyzers_run: Vec<String>,
220 analyzers_skipped: Vec<AnalyzerSkip>,
221 llm_funnel: Option<LlmFunnel>,
222}
223
224impl Engine {
225 pub fn with_builtins() -> Self {
228 Engine {
229 analyzers: crate::analyzer::builtin_analyzers(),
230 policy: crate::policy::Policy::default(),
231 llm: None,
232 ground_llm: None,
233 }
234 }
235
236 pub fn empty() -> Self {
238 Engine {
239 analyzers: vec![],
240 policy: crate::policy::Policy::default(),
241 llm: None,
242 ground_llm: None,
243 }
244 }
245
246 pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
248 self.policy = policy;
249 self
250 }
251
252 pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
257 self.llm = Some(backend);
258 self
259 }
260
261 pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
265 self.ground_llm = Some(backend);
266 self
267 }
268
269 pub fn policy(&self) -> &crate::policy::Policy {
270 &self.policy
271 }
272
273 pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
275 self.analyzers.push(analyzer);
276 }
277
278 pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
279 &self.analyzers
280 }
281
282 pub fn analyze_only<S: OmsSubstrate>(
291 &self,
292 sub: &S,
293 opts: &RunOptions,
294 overrides: &BTreeMap<String, Map<String, Value>>,
295 now_ms: i64,
296 ) -> Result<Vec<Recommendation>> {
297 let persisted = LoopPersisted::from_value(sub.load_state()?)?;
298 let analysis_watermark = if opts.full_sweep {
299 None
300 } else {
301 persisted.state.watermark_ms
302 };
303 Ok(self
304 .analysis_pass(
305 sub,
306 &persisted,
307 opts,
308 overrides,
309 analysis_watermark,
310 now_ms,
311 &[],
312 )?
313 .survivors)
314 }
315
316 pub fn run<S: OmsSubstrate>(
319 &self,
320 sub: &mut S,
321 opts: &RunOptions,
322 now_ms: i64,
323 ) -> Result<RunResult> {
324 let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
325 let watermark = persisted.state.watermark_ms;
326 let analysis_watermark = if opts.full_sweep { None } else { watermark };
331
332 let (new_grains, new_error_events) = count_new(sub, watermark)?;
333 if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
334 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
335 }
336
337 let outcome_inputs = measure_outcomes(sub, &mut persisted, now_ms)?;
340
341 let AnalysisPass {
342 survivors,
343 proposed,
344 deduped,
345 analyzers_run,
346 analyzers_skipped,
347 llm_funnel,
348 } = self.analysis_pass(
349 &*sub,
350 &persisted,
351 opts,
352 &BTreeMap::new(),
353 analysis_watermark,
354 now_ms,
355 &outcome_inputs,
356 )?;
357
358 let mut stored = 0u64;
361 let mut auto_applied = 0u64;
362 for mut rec in survivors {
363 let spec = rec.to_grain_spec(LOOP_NS)?;
364 let hash = sub.put_grain(&spec)?;
365 rec.hash = hash.clone();
366 let actor = format!("engine:{}", rec.analyzer);
367 let audit = AuditRecord {
368 rec_hash: hash.clone(),
369 from: None,
370 to: RecStatus::Pending,
371 actor: actor.clone(),
372 observer_type: ObserverType::System,
373 because: "analyzer proposed".into(),
374 previous_audit_hash: None,
375 gating: None,
376 at_ms: now_ms,
377 };
378 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
379 persisted
380 .status_index
381 .insert(hash.clone(), RecStatus::Pending);
382 persisted.creators.insert(hash.clone(), actor);
383 if !matches!(rec.origin, Origin::Builtin) {
388 if let Some(trigger) = &opts.triggering_actor {
389 persisted.co_creators.insert(hash.clone(), trigger.clone());
390 }
391 }
392 persisted.audit_heads.insert(hash.clone(), audit_hash);
393 stored += 1;
394
395 if self.can_auto_apply(&*sub, &rec) {
396 self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
397 auto_applied += 1;
398 }
399 }
400
401 persisted.state.last_run_ms = Some(now_ms);
402 persisted.state.watermark_ms = Some(now_ms);
403 sub.store_state(&persisted.to_value()?)?;
404
405 Ok(RunResult {
406 outcome: RunOutcome::Ran,
407 skip_reason: None,
408 new_grains,
409 new_error_events,
410 proposed,
411 deduped,
412 stored,
413 auto_applied,
414 analyzers_run,
415 analyzers_skipped,
416 llm_funnel,
417 })
418 }
419
420 #[allow(clippy::too_many_arguments)]
424 fn analysis_pass<S: OmsSubstrate>(
425 &self,
426 sub: &S,
427 persisted: &LoopPersisted,
428 opts: &RunOptions,
429 external_overrides: &BTreeMap<String, Map<String, Value>>,
430 analysis_watermark: Option<i64>,
431 now_ms: i64,
432 outcome_inputs: &[OutcomeInput],
433 ) -> Result<AnalysisPass> {
434 let existing = existing_dedup_keys(sub, persisted)?;
435 let mut analyzers_run = Vec::new();
436 let mut analyzers_skipped = Vec::new();
437 let mut candidates: Vec<Recommendation> = Vec::new();
438 let caps = sub.capabilities();
439
440 for analyzer in &self.analyzers {
441 let m = analyzer.manifest();
442 let cfg = persisted.config.get(&m.id);
443 let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
444 if !enabled {
445 analyzers_skipped.push(AnalyzerSkip {
446 id: m.id.clone(),
447 reason: "disabled".into(),
448 });
449 continue;
450 }
451 if self.policy.denies(m.family()) {
452 analyzers_skipped.push(AnalyzerSkip {
453 id: m.id.clone(),
454 reason: "denied by host policy".into(),
455 });
456 continue;
457 }
458 if let Some(missing) = missing_capability(m, caps) {
459 analyzers_skipped.push(AnalyzerSkip {
460 id: m.id.clone(),
461 reason: format!("missing capability: {missing}"),
462 });
463 continue;
464 }
465 let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
466 if let Some(extra) = external_overrides.get(&m.id) {
467 for (key, value) in extra {
468 param_overrides.insert(key.clone(), value.clone());
469 }
470 }
471 let params = match m.resolve_params(¶m_overrides) {
472 Ok(p) => p,
473 Err(e) => {
474 analyzers_skipped.push(AnalyzerSkip {
475 id: m.id.clone(),
476 reason: e.to_string(),
477 });
478 continue;
479 }
480 };
481 let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
482 let ns_slice: &[String] = if ns_owned.is_empty() {
483 &opts.namespaces
484 } else {
485 &ns_owned
486 };
487 let reader: &dyn SubstrateRead = sub;
488 let ctx = AnalyzeCtx::new(
489 reader,
490 ¶ms,
491 ns_slice,
492 analysis_watermark,
493 now_ms,
494 outcome_inputs,
495 );
496 match analyzer.analyze(&ctx) {
497 Ok(drafts) => {
498 analyzers_run.push(m.id.clone());
499 for draft in drafts {
500 match stamp(m, ¶ms, draft, now_ms) {
501 Ok(rec) => candidates.push(rec),
502 Err(e) => analyzers_skipped.push(AnalyzerSkip {
503 id: m.id.clone(),
504 reason: e.to_string(),
505 }),
506 }
507 }
508 }
509 Err(e) => analyzers_skipped.push(AnalyzerSkip {
510 id: m.id.clone(),
511 reason: e.to_string(),
512 }),
513 }
514 }
515
516 let mut funnel = LlmFunnel::default();
517 if self.llm.is_some() {
518 candidates.extend(self.discover(
519 sub,
520 &candidates,
521 analysis_watermark,
522 &opts.namespaces,
523 now_ms,
524 &mut funnel,
525 ));
526 }
527
528 let proposed = candidates.len() as u64;
529 let mut seen = BTreeSet::new();
530 let mut survivors = Vec::new();
531 for candidate in candidates {
532 let family = crate::manifest::analyzer_family(&candidate.analyzer);
533 let floor = [
534 severity_floor_for(persisted, &candidate.analyzer),
535 self.policy.severity_floor(family),
536 ]
537 .into_iter()
538 .flatten()
539 .max();
540 if floor.is_some_and(|floor| candidate.severity < floor) {
541 continue;
542 }
543 if !seen.insert(candidate.dedup_key.clone()) {
544 continue;
545 }
546 if existing.contains(&candidate.dedup_key) {
547 continue;
548 }
549 if persisted
550 .cooldowns
551 .get(&candidate.dedup_key)
552 .is_some_and(|until| now_ms < *until)
553 {
554 continue;
555 }
556 survivors.push(candidate);
557 }
558 let deduped = proposed - survivors.len() as u64;
559 if self.llm.is_some() {
560 self.enrich(&mut survivors);
561 }
562 Ok(AnalysisPass {
563 survivors,
564 proposed,
565 deduped,
566 analyzers_run,
567 analyzers_skipped,
568 llm_funnel: self.llm.is_some().then_some(funnel),
569 })
570 }
571
572 fn discover<S: OmsSubstrate>(
579 &self,
580 sub: &S,
581 candidates: &[Recommendation],
582 watermark: Option<i64>,
583 namespaces: &[String],
584 now_ms: i64,
585 funnel: &mut LlmFunnel,
586 ) -> Vec<Recommendation> {
587 let Some(llm) = &self.llm else {
588 return Vec::new();
589 };
590 let findings: Vec<crate::llm::FindingBrief> = candidates
591 .iter()
592 .take(32)
593 .map(|c| crate::llm::FindingBrief {
594 analyzer: c.analyzer.clone(),
595 summary: c.summary.render(),
596 target: c.target_ref.clone(),
597 severity: c.severity.as_str().to_string(),
598 })
599 .collect();
600 let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
606 let mut bundle: BTreeSet<String> = BTreeSet::new();
607 let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
608 'cited: for c in candidates {
609 for h in &c.evidence {
610 if evidence.len() >= CITED_SEED_CAP {
619 break 'cited;
620 }
621 if !bundle.contains(h) {
622 if let Ok(Some(g)) = sub.grain(h) {
623 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
624 }
625 }
626 }
627 }
628 let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
629 vec![None]
630 } else {
631 namespaces.iter().map(|n| Some(n.as_str())).collect()
632 };
633 let opts = ReadOpts { live_only: true, since_ms: watermark };
634 let mut tool_seeded = 0usize;
646 'tools: for ns in &scan_ns {
647 if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
648 for g in recent {
649 if tool_seeded >= TOOL_SEED_CAP
650 || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
651 {
652 break 'tools;
653 }
654 if !g.is_error() {
655 continue;
656 }
657 let before = evidence.len();
658 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
659 if evidence.len() > before {
660 tool_seeded += 1;
661 }
662 }
663 }
664 }
665 'notes: for ns in &scan_ns {
675 if let Ok(recent) =
676 sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
677 {
678 for g in recent {
679 if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
680 break 'notes;
681 }
682 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
683 }
684 }
685 }
686 'seed: for gt in [
687 crate::model::grain_type::FACT,
688 crate::model::grain_type::OBSERVATION,
689 ] {
690 for ns in &scan_ns {
691 if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
692 for g in recent {
693 if evidence.len() >= EVIDENCE_CAP {
694 break 'seed;
695 }
696 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
697 }
698 }
699 }
700 }
701 funnel.evidence = evidence.len() as u64;
702 if evidence.is_empty() {
703 return Vec::new(); }
705 let (approved, rejected) = self.llm_history(sub);
710 let request = crate::llm::LlmRequest {
711 loop_proto: 1,
712 op: "discover",
713 instructions: DISCOVER_INSTRUCTIONS,
714 findings: findings.clone(),
715 evidence: evidence.clone(),
716 rejected,
717 approved,
718 };
719 let Ok(body) = serde_json::to_string(&request) else {
720 return Vec::new();
721 };
722 let raw = match llm.complete(&body) {
723 Ok(r) => r,
724 Err(_) => return Vec::new(), };
726 let caps = sub.capabilities();
730 let mut validated: Vec<ValidatedDraft> = Vec::new();
731 let drafts: Vec<_> = crate::llm::parse_discover(&raw)
732 .recommendations
733 .into_iter()
734 .take(crate::llm::MAX_LLM_DRAFTS)
735 .collect();
736 funnel.proposed = drafts.len() as u64;
737 for d in drafts {
738 let cited: Vec<String> =
739 d.evidence.iter().filter(|h| bundle.contains(*h)).cloned().collect();
740 if cited.is_empty() {
741 funnel.dropped_uncited += 1;
742 continue; }
744 let Ok(target) = TargetRef::parse(&d.target) else {
745 funnel.dropped_target += 1;
746 continue;
747 };
748 let tc = target.target_class();
749 if !matches!(tc, "memory" | "query" | "code") {
755 funnel.dropped_target += 1;
756 continue;
757 }
758 let resolved = resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps);
761 if tc == "code" && resolved.is_none() {
766 funnel.dropped_target += 1;
767 continue;
768 }
769 validated.push(ValidatedDraft {
770 draft: d,
771 target_ref: target.as_string(),
772 cited,
773 resolved,
774 });
775 }
776 funnel.cited = validated.len() as u64;
777 if validated.is_empty() {
778 return Vec::new();
779 }
780 let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
787 self.verify_drafts(&**llm, ground, validated, &evidence, now_ms, funnel)
788 }
789
790 fn verify_drafts(
797 &self,
798 llm: &dyn crate::llm::LlmBackend,
799 ground: &dyn crate::llm::LlmBackend,
800 validated: Vec<ValidatedDraft>,
801 evidence: &[crate::llm::EvidenceItem],
802 now_ms: i64,
803 funnel: &mut LlmFunnel,
804 ) -> Vec<Recommendation> {
805 use crate::llm::*;
806 let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
807 evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
808 let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
809 cited
810 .iter()
811 .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
812 .collect()
813 };
814
815 let claims: Vec<GroundItem> = validated
819 .iter()
820 .enumerate()
821 .map(|(i, v)| GroundItem {
822 id: i,
823 claim: claim_text(&v.draft, v.resolved.as_ref()),
824 evidence: ev_for(&v.cited),
825 })
826 .collect();
827 let ground_req = GroundRequest {
828 loop_proto: 1,
829 op: "ground",
830 instructions: GROUND_INSTRUCTIONS,
831 claims,
832 };
833 let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
839 .ok()
840 .and_then(|b| ground.complete(&b).ok())
841 {
842 Some(raw) => {
843 let parsed = parse_ground(&raw);
844 funnel.ground_verdicts = parsed.results.len() as u64;
845 parsed
846 .results
847 .into_iter()
848 .filter(|r| r.supported)
849 .map(|r| r.id)
850 .collect()
851 }
852 None => {
853 funnel.ground_call_failed = true;
854 return Vec::new();
855 }
856 };
857 funnel.grounded = grounded.len() as u64;
858 if grounded.is_empty() {
859 return Vec::new();
860 }
861
862 let items: Vec<VerifyItem> = validated
868 .iter()
869 .enumerate()
870 .filter(|(i, _)| grounded.contains(i))
871 .map(|(i, v)| VerifyItem {
872 id: i,
873 summary: claim_text(&v.draft, v.resolved.as_ref()),
875 target: v.target_ref.clone(),
876 evidence: ev_for(&v.cited),
877 })
878 .collect();
879 let verify_req = VerifyRequest {
880 loop_proto: 1,
881 op: "verify",
882 instructions: VERIFY_INSTRUCTIONS,
883 findings: items,
884 };
885 let verdicts: std::collections::BTreeMap<usize, f64> =
886 match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
887 Some(raw) => parse_verify(&raw)
888 .results
889 .into_iter()
890 .filter(|r| r.keep)
891 .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
892 .collect(),
893 None => return Vec::new(),
894 };
895
896 funnel.kept = verdicts.len() as u64;
897 let mut out = Vec::new();
901 for (i, v) in validated.into_iter().enumerate() {
902 if let Some(&conf) = verdicts.get(&i) {
903 if conf >= MIN_LLM_CONFIDENCE {
904 out.push(stamp_llm(
905 llm.model(),
906 &v.draft,
907 v.target_ref,
908 v.cited,
909 v.resolved,
910 conf,
911 now_ms,
912 ));
913 }
914 }
915 }
916 funnel.stored = out.len() as u64;
917 out
918 }
919
920 fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
925 const MAX: usize = 20;
926 let Ok(mut recs) = self.recommendations(sub, None) else {
927 return (Vec::new(), Vec::new());
928 };
929 recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
930 recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
931 let mut approved = Vec::new();
932 let mut rejected = Vec::new();
933 for r in &recs {
934 match r.status {
935 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
936 if approved.len() < MAX =>
937 {
938 approved.push(r.summary.render());
939 }
940 RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
941 _ => {}
942 }
943 }
944 (approved, rejected)
945 }
946
947 fn enrich(&self, survivors: &mut [Recommendation]) {
952 let Some(llm) = &self.llm else {
953 return;
954 };
955 if survivors.is_empty() {
956 return;
957 }
958 let findings: Vec<crate::llm::FindingBrief> = survivors
959 .iter()
960 .map(|r| crate::llm::FindingBrief {
961 analyzer: r.analyzer.clone(),
962 summary: r.summary.render(),
963 target: r.target_ref.clone(),
964 severity: r.severity.as_str().to_string(),
965 })
966 .collect();
967 let request = crate::llm::LlmRequest {
968 loop_proto: 1,
969 op: "enrich",
970 instructions: ENRICH_INSTRUCTIONS,
971 findings,
972 evidence: Vec::new(),
973 rejected: Vec::new(),
974 approved: Vec::new(),
975 };
976 let Ok(body) = serde_json::to_string(&request) else {
977 return;
978 };
979 let raw = match llm.complete(&body) {
980 Ok(r) => r,
981 Err(_) => return,
982 };
983 for note in crate::llm::parse_enrich(&raw).notes {
984 if note.guidance.trim().is_empty() {
985 continue;
986 }
987 if let Some(r) = survivors
988 .iter_mut()
989 .find(|r| r.target_ref == note.target && r.guidance.is_none())
990 {
991 r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
992 }
993 }
994 }
995
996 fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
1005 if !rec.origin.auto_apply_eligible() || rec.destructive {
1006 return false;
1007 }
1008 let manifest_ok = self
1012 .analyzers
1013 .iter()
1014 .map(|a| a.manifest())
1015 .find(|m| m.id == rec.analyzer)
1016 .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
1017 if !manifest_ok {
1018 return false;
1019 }
1020 let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1021 return false;
1022 };
1023 let family = crate::manifest::analyzer_family(&rec.analyzer);
1024 if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1025 return false;
1026 }
1027 match &rec.proposal {
1032 Proposal::Cal { cal } => cal
1033 .lines()
1034 .map(str::trim)
1035 .filter(|l| !l.is_empty())
1036 .all(|l| supersede_is_value_identical(sub, l)),
1037 _ => false,
1038 }
1039 }
1040
1041 fn auto_apply<S: OmsSubstrate>(
1044 &self,
1045 sub: &mut S,
1046 p: &mut LoopPersisted,
1047 rec: &Recommendation,
1048 now_ms: i64,
1049 ) -> Result<()> {
1050 let mut created = Vec::new();
1051 if let Proposal::Cal { cal } = &rec.proposal {
1052 if cal.lines().map(str::trim).any(is_definition_statement) {
1057 return Err(Error::InvalidProposal(
1058 "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1059 auto-applied: it changes what every future context contains, so it \
1060 requires a human APPROVE + APPLY with BECAUSE"
1061 .into(),
1062 ));
1063 }
1064 for r in sub.execute_cal(cal)? {
1065 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1066 created.push(h.to_string());
1067 }
1068 }
1069 }
1070 let applied = AppliedRecord {
1071 applied_at_ms: now_ms,
1072 target_ref: rec.target_ref.clone(),
1073 rollbackable: rec.rollbackable,
1074 created_hashes: created,
1075 inverse_cal: None,
1076 metric: rec.metric.clone(),
1077 };
1078 let prev = p.audit_heads.get(&rec.hash).cloned();
1079 let audit = AuditRecord {
1080 rec_hash: rec.hash.clone(),
1081 from: Some(RecStatus::Pending),
1082 to: RecStatus::Applied,
1083 actor: "policy:auto".into(),
1084 observer_type: ObserverType::Policy,
1085 because: "auto-applied per host policy".into(),
1086 previous_audit_hash: prev,
1087 gating: None,
1088 at_ms: now_ms,
1089 };
1090 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1091 p.audit_heads.insert(rec.hash.clone(), audit_hash);
1092 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1093 p.applied.insert(rec.hash.clone(), applied);
1094 Ok(())
1095 }
1096
1097 #[allow(clippy::too_many_arguments)]
1100 pub fn review<S: OmsSubstrate>(
1101 &self,
1102 sub: &mut S,
1103 rec_hash: &str,
1104 decision: Decision,
1105 actor: &str,
1106 observer: ObserverType,
1107 scopes: &ScopeSet,
1108 because: &str,
1109 now_ms: i64,
1110 ) -> Result<()> {
1111 if !scopes.has(Scope::Review) {
1112 return Err(Error::ScopeDenied("review".into()));
1113 }
1114 let because = validate_because(because)?;
1115 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1116 let status = *p
1117 .status_index
1118 .get(rec_hash)
1119 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1120 let to = match decision {
1121 Decision::Approve => RecStatus::Approved,
1122 Decision::Reject => RecStatus::Rejected,
1123 };
1124 if !status.can_transition_to(to, false) {
1125 return Err(Error::LifecycleViolation(format!(
1126 "{} -> {}",
1127 status.as_str(),
1128 to.as_str()
1129 )));
1130 }
1131 if to == RecStatus::Approved {
1132 if let Some(creator) = p.creators.get(rec_hash) {
1133 if creator == actor {
1134 return Err(Error::SelfApproval(format!(
1135 "{actor} created this recommendation"
1136 )));
1137 }
1138 }
1139 if let Some(trigger) = p.co_creators.get(rec_hash) {
1140 if trigger == actor {
1141 return Err(Error::SelfApproval(format!(
1142 "{actor} triggered the run that authored this recommendation"
1143 )));
1144 }
1145 }
1146 }
1147 let prev = p.audit_heads.get(rec_hash).cloned();
1148 let audit = AuditRecord {
1149 rec_hash: rec_hash.into(),
1150 from: Some(status),
1151 to,
1152 actor: actor.into(),
1153 observer_type: observer,
1154 because,
1155 previous_audit_hash: prev,
1156 gating: None,
1157 at_ms: now_ms,
1158 };
1159 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1160 p.audit_heads.insert(rec_hash.into(), audit_hash);
1161 p.status_index.insert(rec_hash.into(), to);
1162 if to == RecStatus::Rejected {
1163 if let Ok(rec) = load_rec(sub, rec_hash) {
1164 const BASE_MS: i64 = 7 * 86_400_000;
1169 const CAP_MS: i64 = 90 * 86_400_000;
1170 let strikes = p.cooldown_strikes.entry(rec.dedup_key.clone()).or_insert(0);
1171 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
1172 *strikes = strikes.saturating_add(1);
1173 p.cooldowns.insert(rec.dedup_key, now_ms + interval);
1174 }
1175 }
1176 sub.store_state(&p.to_value()?)?;
1177 Ok(())
1178 }
1179
1180 pub fn preflight_apply<S: OmsSubstrate>(
1197 &self,
1198 sub: &S,
1199 rec_hash: &str,
1200 scopes: &ScopeSet,
1201 allow_destructive: bool,
1202 has_gating: bool,
1203 ) -> Result<()> {
1204 if !scopes.has(Scope::Apply) {
1205 return Err(Error::ScopeDenied("apply".into()));
1206 }
1207 let rec = load_rec(sub, rec_hash)?;
1208 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1209 return Err(Error::DestructiveGated(
1210 "destructive apply requires admin scope + allow_destructive".into(),
1211 ));
1212 }
1213 ensure_executable(rec.action_kind, &rec.proposal)?;
1214 if requires_gating(rec.action_kind) && !has_gating {
1215 return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1216 }
1217 Ok(())
1218 }
1219
1220 #[allow(clippy::too_many_arguments)]
1224 pub fn apply<S: OmsSubstrate>(
1225 &self,
1226 sub: &mut S,
1227 rec_hash: &str,
1228 actor: &str,
1229 observer: ObserverType,
1230 scopes: &ScopeSet,
1231 because: &str,
1232 allow_destructive: bool,
1233 now_ms: i64,
1234 ) -> Result<AppliedRecord> {
1235 self.apply_inner(
1236 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1237 )
1238 }
1239
1240 pub fn gating_evidence<S: OmsSubstrate>(
1247 &self,
1248 sub: &S,
1249 rec_hash: &str,
1250 run_id: &str,
1251 ) -> Result<crate::recommendation::GatingEvidence> {
1252 let rec = self
1253 .recommendations(sub, None)?
1254 .into_iter()
1255 .find(|r| r.hash == rec_hash)
1256 .ok_or_else(|| {
1257 Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1258 })?;
1259 let pin = rec.evalset_hash.ok_or_else(|| {
1260 Error::InvalidProposal(
1261 "this recommendation pins no evalset — a gating run applies only \
1262 to code and adapter revisions"
1263 .into(),
1264 )
1265 })?;
1266 match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1271 Some(run) => Ok(crate::recommendation::GatingEvidence {
1272 evalset_hash: pin,
1273 run_id: run.run_id,
1274 passed: run.passed,
1275 failed: run.failed,
1276 }),
1277 None => Err(Error::InvalidProposal(format!(
1278 "no recorded gate run '{run_id}' for evalset {pin} — run \
1279 `areev eval run --evalset {pin} ...` first"
1280 ))),
1281 }
1282 }
1283
1284 #[allow(clippy::too_many_arguments)]
1288 pub fn apply_gated<S: OmsSubstrate>(
1289 &self,
1290 sub: &mut S,
1291 rec_hash: &str,
1292 actor: &str,
1293 observer: ObserverType,
1294 scopes: &ScopeSet,
1295 because: &str,
1296 allow_destructive: bool,
1297 gating: &crate::recommendation::GatingEvidence,
1298 now_ms: i64,
1299 ) -> Result<AppliedRecord> {
1300 self.apply_inner(
1301 sub,
1302 rec_hash,
1303 actor,
1304 observer,
1305 scopes,
1306 because,
1307 allow_destructive,
1308 Some(gating),
1309 now_ms,
1310 )
1311 }
1312
1313 #[allow(clippy::too_many_arguments)]
1314 fn apply_inner<S: OmsSubstrate>(
1315 &self,
1316 sub: &mut S,
1317 rec_hash: &str,
1318 actor: &str,
1319 observer: ObserverType,
1320 scopes: &ScopeSet,
1321 because: &str,
1322 allow_destructive: bool,
1323 gating: Option<&crate::recommendation::GatingEvidence>,
1324 now_ms: i64,
1325 ) -> Result<AppliedRecord> {
1326 if !scopes.has(Scope::Apply) {
1327 return Err(Error::ScopeDenied("apply".into()));
1328 }
1329 let because = validate_because(because)?;
1330 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1331 let status = *p
1332 .status_index
1333 .get(rec_hash)
1334 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1335 if !status.can_transition_to(RecStatus::Applied, false) {
1336 return Err(Error::LifecycleViolation(format!(
1337 "{} -> applied (approve first)",
1338 status.as_str()
1339 )));
1340 }
1341 let rec = load_rec(sub, rec_hash)?;
1342 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1343 return Err(Error::DestructiveGated(
1344 "destructive apply requires admin scope + allow_destructive".into(),
1345 ));
1346 }
1347 if requires_gating(rec.action_kind) {
1352 let g = gating
1353 .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1354 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1355 if g.evalset_hash != pin {
1356 return Err(Error::InvalidProposal(format!(
1357 "gating ran evalset {} but the recommendation is pinned \
1358 to {pin} (Rule E1)",
1359 g.evalset_hash
1360 )));
1361 }
1362 match sub.grain(pin)? {
1363 Some(evalset) if evalset.is_live() => {}
1364 Some(_) => {
1365 return Err(Error::InvalidProposal(
1366 "the pinned evalset was superseded after gating — \
1367 the recommendation must re-gate (Rule E1)"
1368 .into(),
1369 ))
1370 }
1371 None => {
1372 return Err(Error::InvalidProposal(format!(
1373 "pinned evalset {pin} not found in the substrate"
1374 )))
1375 }
1376 }
1377 if g.failed > 0 {
1378 return Err(Error::InvalidProposal(format!(
1379 "the gating run failed {}/{} cases — a failing gate \
1380 admits nothing",
1381 g.failed,
1382 g.passed + g.failed
1383 )));
1384 }
1385 }
1386
1387 let mut created = Vec::new();
1389 let mut inverse_cal: Option<String> = None;
1393 match &rec.proposal {
1394 Proposal::Cal { cal } => {
1395 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1396 if !is_definition_statement(line) {
1397 continue;
1398 }
1399 match sub.definition_inverse(line)? {
1400 Some(inv) => inverse_cal = Some(inv),
1401 None => {
1402 return Err(Error::InvalidProposal(format!(
1403 "this substrate cannot record a rollback inverse for {line:?}; \
1404 a definition rewrite that ROLLBACK could not undo is refused \
1405 rather than applied"
1406 )))
1407 }
1408 }
1409 }
1410 let rows = sub.execute_cal(cal)?;
1411 for r in rows {
1412 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1413 created.push(h.to_string());
1414 }
1415 }
1416 }
1417 Proposal::Edit { .. } => {
1420 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1421 }
1422 Proposal::Data { data } if requires_gating(rec.action_kind) => {
1430 let relation = if rec.action_kind == ActionKind::AdapterRevision {
1431 "mg:adapter_promotion"
1432 } else {
1433 "mg:code_promotion"
1434 };
1435 let mut promoted = data.clone();
1443 if let Some(Value::String(src)) = promoted.remove("source") {
1444 let address = sub.put_blob(src.as_bytes())?;
1445 promoted.insert("code_address".into(), Value::from(address));
1446 }
1447 let mut spec = crate::substrate::GrainSpec::new(
1448 crate::model::grain_type::FACT,
1449 LOOP_NS,
1450 )
1451 .with_field("subject", rec.target_ref.clone())
1452 .with_field("relation", relation)
1453 .with_field(
1454 "object",
1455 serde_json::to_string(&promoted)
1456 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1457 )
1458 .with_field("rec_hash", rec_hash.to_string());
1459 if let Some(g) = gating {
1460 spec = spec
1461 .with_field("gating_evalset", g.evalset_hash.clone())
1462 .with_field("gating_run_id", g.run_id.clone());
1463 }
1464 created.push(sub.put_grain(&spec)?);
1465 }
1466 Proposal::Data { data } => {
1467 let revert_of = data
1472 .get("revert_of")
1473 .and_then(Value::as_str)
1474 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1475 self.rollback(
1476 sub,
1477 revert_of,
1478 actor,
1479 observer,
1480 scopes,
1481 &because,
1482 now_ms,
1483 )?;
1484 p = LoopPersisted::from_value(sub.load_state()?)?;
1487 }
1488 }
1489
1490 let applied = AppliedRecord {
1491 applied_at_ms: now_ms,
1492 target_ref: rec.target_ref.clone(),
1493 rollbackable: rec.rollbackable,
1494 created_hashes: created,
1495 inverse_cal,
1496 metric: rec.metric.clone(),
1497 };
1498 let prev = p.audit_heads.get(rec_hash).cloned();
1499 let audit = AuditRecord {
1500 rec_hash: rec_hash.into(),
1501 from: Some(status),
1502 to: RecStatus::Applied,
1503 actor: actor.into(),
1504 observer_type: observer,
1505 because,
1506 previous_audit_hash: prev,
1507 gating: gating.cloned(),
1508 at_ms: now_ms,
1509 };
1510 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1511 p.audit_heads.insert(rec_hash.into(), audit_hash);
1512 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1513 p.applied.insert(rec_hash.into(), applied.clone());
1514 sub.store_state(&p.to_value()?)?;
1515 Ok(applied)
1516 }
1517
1518 #[allow(clippy::too_many_arguments)]
1521 pub fn rollback<S: OmsSubstrate>(
1522 &self,
1523 sub: &mut S,
1524 rec_hash: &str,
1525 actor: &str,
1526 observer: ObserverType,
1527 scopes: &ScopeSet,
1528 because: &str,
1529 now_ms: i64,
1530 ) -> Result<()> {
1531 if !scopes.has(Scope::Apply) {
1532 return Err(Error::ScopeDenied("apply".into()));
1533 }
1534 let because = validate_because(because)?;
1535 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1536 let status = *p
1537 .status_index
1538 .get(rec_hash)
1539 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1540 if !status.can_transition_to(RecStatus::RolledBack, false) {
1541 return Err(Error::LifecycleViolation(format!(
1542 "{} -> rolled_back",
1543 status.as_str()
1544 )));
1545 }
1546 let applied = p
1547 .applied
1548 .get(rec_hash)
1549 .cloned()
1550 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1551 if !applied.rollbackable {
1552 return Err(Error::LifecycleViolation(
1553 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1554 ));
1555 }
1556 for h in &applied.created_hashes {
1557 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1558 }
1559 if let Some(inverse) = &applied.inverse_cal {
1566 sub.execute_cal(inverse)?;
1567 }
1568 let prev = p.audit_heads.get(rec_hash).cloned();
1569 let audit = AuditRecord {
1570 rec_hash: rec_hash.into(),
1571 from: Some(status),
1572 to: RecStatus::RolledBack,
1573 actor: actor.into(),
1574 observer_type: observer,
1575 because,
1576 previous_audit_hash: prev,
1577 gating: None,
1578 at_ms: now_ms,
1579 };
1580 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1581 p.audit_heads.insert(rec_hash.into(), audit_hash);
1582 p.status_index
1583 .insert(rec_hash.into(), RecStatus::RolledBack);
1584 sub.store_state(&p.to_value()?)?;
1585 Ok(())
1586 }
1587
1588 pub fn recommendations<S: OmsSubstrate>(
1593 &self,
1594 sub: &S,
1595 status_filter: Option<RecStatus>,
1596 ) -> Result<Vec<Recommendation>> {
1597 let p = LoopPersisted::from_value(sub.load_state()?)?;
1598 let grains = sub.grains_of_type(
1599 crate::model::grain_type::RECOMMENDATION,
1600 Some(LOOP_NS),
1601 ReadOpts {
1602 live_only: false,
1603 since_ms: None,
1604 },
1605 )?;
1606 let mut out = Vec::new();
1607 for g in grains {
1608 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1609 rec.status = p
1610 .status_index
1611 .get(&g.hash)
1612 .copied()
1613 .unwrap_or(RecStatus::Pending);
1614 if let Some(f) = status_filter {
1615 if rec.status != f {
1616 continue;
1617 }
1618 }
1619 out.push(rec);
1620 }
1621 out.sort_by(|a, b| {
1630 b.severity
1631 .cmp(&a.severity)
1632 .then(a.created_at_ms.cmp(&b.created_at_ms))
1633 .then(a.dedup_key.cmp(&b.dedup_key))
1634 .then(a.hash.cmp(&b.hash))
1635 });
1636 Ok(out)
1637 }
1638
1639 pub fn analyzer_settings<S: OmsSubstrate>(
1642 &self,
1643 sub: &S,
1644 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1645 let p = LoopPersisted::from_value(sub.load_state()?)?;
1646 Ok(self
1647 .analyzers
1648 .iter()
1649 .map(|a| {
1650 let m = a.manifest();
1651 let cfg = p.config.get(&m.id);
1652 crate::config::AnalyzerSetting {
1653 id: m.id.clone(),
1654 title: m.title.clone(),
1655 description: m.description.clone(),
1656 tier: format!("{:?}", m.tier),
1657 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1658 default_on: m.default_on,
1659 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1660 severity_floor: cfg
1661 .and_then(|c| c.severity_floor)
1662 .map(|s| s.as_str().to_string()),
1663 }
1664 })
1665 .collect())
1666 }
1667
1668 pub fn set_analyzer_config<S: OmsSubstrate>(
1674 &self,
1675 sub: &mut S,
1676 analyzer_id: &str,
1677 update: crate::config::AnalyzerConfigUpdate,
1678 scopes: &ScopeSet,
1679 ) -> Result<crate::config::AnalyzerConfig> {
1680 if !scopes.has(Scope::Admin) {
1681 return Err(Error::ScopeDenied("admin".into()));
1682 }
1683 let manifest = self
1684 .analyzers
1685 .iter()
1686 .map(|a| a.manifest())
1687 .find(|m| m.id == analyzer_id)
1688 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1689 if let Some(params) = &update.params {
1691 manifest.resolve_params(params)?;
1692 }
1693 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1694 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1695 if let Some(enabled) = update.enabled {
1696 cfg.enabled = Some(enabled);
1697 }
1698 if update.clear_floor {
1699 cfg.severity_floor = None;
1700 } else if let Some(floor) = update.severity_floor {
1701 cfg.severity_floor = Some(floor);
1702 }
1703 if let Some(params) = update.params {
1704 cfg.params = params;
1705 }
1706 if let Some(ns) = update.namespaces {
1707 cfg.namespaces = ns;
1708 }
1709 let stored = cfg.clone();
1710 sub.store_state(&p.to_value()?)?;
1711 Ok(stored)
1712 }
1713
1714 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1717 let p = LoopPersisted::from_value(sub.load_state()?)?;
1718 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1719 out.sort_by(|a, b| {
1723 a.measured_at_ms
1724 .cmp(&b.measured_at_ms)
1725 .then(a.horizon_ms.cmp(&b.horizon_ms))
1726 .then(a.metric.cmp(&b.metric))
1727 .then(a.rec_hash.cmp(&b.rec_hash))
1728 });
1729 Ok(out)
1730 }
1731
1732 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1736 let p = LoopPersisted::from_value(sub.load_state()?)?;
1737 let (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
1738 let recs = self.recommendations(sub, None)?;
1739 let mut pending = 0;
1740 let mut applied = 0;
1741 for r in &recs {
1742 match r.status {
1743 RecStatus::Pending => pending += 1,
1744 RecStatus::Applied => applied += 1,
1745 _ => {}
1746 }
1747 }
1748 let stale = match p.state.last_run_ms {
1750 None => true,
1751 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1752 };
1753 Ok(Health {
1754 last_run_ms: p.state.last_run_ms,
1755 grains_since_run,
1756 error_events_since_run,
1757 pending,
1758 applied,
1759 total: recs.len() as u64,
1760 stale,
1761 })
1762 }
1763
1764 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1769 let recs = self.recommendations(sub, None)?;
1770 let mut m = LlmMetrics::default();
1771 for r in &recs {
1772 if !matches!(r.origin, Origin::Llm { .. }) {
1773 continue;
1774 }
1775 m.proposed += 1;
1776 match r.status {
1777 RecStatus::Pending => m.pending += 1,
1778 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1779 RecStatus::Rejected => m.rejected += 1,
1780 RecStatus::Expired => {}
1781 }
1782 }
1783 let decided = m.approved + m.rejected;
1784 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1785 Ok(m)
1786 }
1787}
1788
1789#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1791pub struct Health {
1792 #[serde(skip_serializing_if = "Option::is_none")]
1793 pub last_run_ms: Option<i64>,
1794 pub grains_since_run: u64,
1795 pub error_events_since_run: u64,
1796 pub pending: u64,
1797 pub applied: u64,
1798 pub total: u64,
1799 pub stale: bool,
1802}
1803
1804#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1806pub struct LlmMetrics {
1807 pub proposed: u64,
1810 pub pending: u64,
1811 pub approved: u64,
1813 pub rejected: u64,
1814 #[serde(skip_serializing_if = "Option::is_none")]
1816 pub approval_rate: Option<f64>,
1817}
1818
1819fn measure_outcomes<S: OmsSubstrate>(
1826 sub: &S,
1827 p: &mut LoopPersisted,
1828 now_ms: i64,
1829) -> Result<Vec<OutcomeInput>> {
1830 let mut due: Vec<(String, crate::config::AppliedRecord, i64)> = Vec::new();
1832 for (h, a) in &p.applied {
1833 if p.status_index.get(h) != Some(&RecStatus::Applied) {
1834 continue;
1835 }
1836 let Some(metric) = &a.metric else { continue };
1837 let done = p.measured.get(h).cloned().unwrap_or_default();
1838 for horizon in metric.horizons() {
1839 if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
1840 due.push((h.clone(), a.clone(), horizon));
1841 }
1842 }
1843 }
1844
1845 let mut out = Vec::new();
1846 for (rec_hash, applied, horizon) in due {
1847 let metric = applied.metric.as_ref().unwrap();
1848 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1849 continue; };
1851 let regressed = crate::recommendation::is_regression(
1852 metric.baseline,
1853 current,
1854 metric.higher_is_better,
1855 );
1856 p.outcomes.entry(rec_hash.clone()).or_default().push(
1857 crate::recommendation::OutcomeResult {
1858 rec_hash: rec_hash.clone(),
1859 metric: metric.metric.clone(),
1860 baseline: metric.baseline,
1861 current,
1862 verdict: if regressed { "regressed" } else { "held" }.into(),
1863 horizon_ms: horizon,
1864 measured_at_ms: now_ms,
1865 },
1866 );
1867 p.measured.entry(rec_hash.clone()).or_default().push(horizon);
1868 if regressed {
1869 out.push(OutcomeInput {
1870 rec_hash,
1871 target_ref: applied.target_ref.clone(),
1872 metric: metric.metric.clone(),
1873 baseline: metric.baseline,
1874 current,
1875 unit: metric.unit.clone(),
1876 higher_is_better: metric.higher_is_better,
1877 });
1878 }
1879 }
1880 Ok(out)
1881}
1882
1883pub(crate) fn measure_metric<S: SubstrateRead>(
1885 sub: &S,
1886 metric: &crate::recommendation::MetricSnapshot,
1887 since_ms: i64,
1888) -> Result<Option<f64>> {
1889 match metric.metric.as_str() {
1890 "tool_error_recurrence" => {
1895 let Some(tool) = &metric.subject else { return Ok(None) };
1896 let tools = sub.grains_of_type(
1897 crate::model::grain_type::TOOL,
1898 None,
1899 ReadOpts { live_only: true, since_ms: Some(since_ms) },
1900 )?;
1901 let n = tools
1902 .iter()
1903 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
1904 .filter(|t| {
1905 metric.relation.as_deref().is_none_or(|sig| {
1908 crate::analyzers::tool_failure::normalize_signature(
1909 t.tool_content().unwrap_or(""),
1910 ) == sig
1911 })
1912 })
1913 .count();
1914 Ok(Some(n as f64))
1915 }
1916 "contradiction_recurrence" => {
1920 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
1921 return Ok(None);
1922 };
1923 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
1924 let distinct: BTreeSet<String> = facts
1925 .iter()
1926 .filter(|f| {
1927 f.fact_relation()
1928 .is_some_and(|r| normalize_ident(r) == *relation)
1929 })
1930 .filter_map(|f| f.fact_object().map(normalize_ident))
1931 .collect();
1932 Ok(Some(distinct.len().saturating_sub(1) as f64))
1933 }
1934 m if m.starts_with("evalset:") => {
1947 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
1948 return Ok(None);
1949 };
1950 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
1951 return Ok(None);
1952 };
1953 Ok(match field {
1957 "failed" => Some(run.failed as f64),
1958 "passed" => Some(run.passed as f64),
1959 "total" => Some(run.total() as f64),
1960 "error_rate" => match run.total() {
1961 0 => None, t => Some(run.failed as f64 / t as f64),
1963 },
1964 other => run.field(other),
1965 })
1966 }
1967 _ => Ok(None),
1968 }
1969}
1970
1971fn scoped_live_facts<S: SubstrateRead>(
1974 sub: &S,
1975 namespace: Option<&str>,
1976 subject: &str,
1977) -> Result<Vec<GrainRecord>> {
1978 let facts = sub.grains_of_type(
1979 crate::model::grain_type::FACT,
1980 None,
1981 ReadOpts { live_only: true, since_ms: None },
1982 )?;
1983 Ok(facts
1984 .into_iter()
1985 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
1986 .filter(|f| {
1987 f.fact_subject()
1988 .is_some_and(|s| normalize_ident(s) == subject)
1989 })
1990 .collect())
1991}
1992
1993fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
2005 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
2006 return false;
2007 };
2008 if fields.is_empty() {
2009 return false;
2010 }
2011 let Ok(Some(grain)) = sub.grain(&target) else {
2012 return false;
2013 };
2014 if grain.valid_to_ms.is_some() {
2024 return false;
2025 }
2026 fields.iter().all(|(k, v)| {
2028 if k == "namespace" {
2029 return v
2030 .as_str()
2031 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2032 }
2033 match (v, grain.fields.get(k)) {
2034 (Value::String(a), Some(Value::String(b))) => {
2035 normalize_ident(a) == normalize_ident(b)
2036 }
2037 (a, Some(b)) => a == b,
2038 (_, None) => false,
2039 }
2040 })
2041}
2042
2043const EVIDENCE_CAP: usize = 64;
2053const CITED_SEED_CAP: usize = 24;
2054const TOOL_SEED_CAP: usize = 16;
2055const NOTE_SEED_CAP: usize = 8;
2060const LENS_RESERVE: usize = 24;
2061
2062const MIN_LLM_CONFIDENCE: f64 = 0.75;
2065
2066const DISCOVER_INSTRUCTIONS: &str = "You review an agent's memory for quality. \
2071Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2072findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2073stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2074recurring cost or hand-off the agent's own setup could remove). \
2075The deterministic findings already cover what the ERROR TEXT says; restating \
2076one of them earns nothing. The evidence may also contain OUTCOME records — a \
2077run's observable shape together with whether it was accepted or rejected. A \
2078problem that raised no error at all is exactly the kind the deterministic \
2079checks cannot see, so compare the rejected outcomes against the accepted \
2080ones: a feature they share and the accepted ones lack is a candidate rule. \
2081Require at least two rejected outcomes before proposing one — a single \
2082rejection is an anecdote, not a pattern. \
2083SCORING: propose a finding ONLY if you \
2084are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2085useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2086earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2087when there is nothing worth flagging. The 'approved' and 'rejected' lists, when \
2088present, show findings this reviewer recently accepted or rejected — prefer the \
2089kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2090or more evidence hashes from the bundle, name a 'target', and include your \
2091confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
2092\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
2093\"confidence\":0.0,\"proposal\":{...}}]}. \
2094OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2095you are not asking to change anything. Include it ONLY when the evidence \
2096supports a specific change, choosing exactly one kind: \
2097(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2098\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
2099an action the agent itself takes on the next occasion. Either ADD an action \
2100it is failing to take ('Record the vendor name and the amount on every \
2101invoice, not just the date') or ORDER one that goes wrong ('Refund a \
2102subscription before cancelling it; refunds on cancelled subscriptions are \
2103refused'). Name the action, not a check on it: 'validate', 'verify' and \
2104'ensure ... is correct' describe a review step the agent has no way to \
2105perform, and such a rule changes nothing even once applied. \
2106(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2107entity target — a durable fact the agent keeps having to be told (an alias, a \
2108settled default, a preference). 'relation' is a short identifier (letters, \
2109digits, _ - . :), not a sentence. \
2110(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2111\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2112assembles the agent's context, when the evidence shows it retrieves the wrong \
2113things. Give the FULL new body; it replaces the old one. \
2114(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2115with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2116workflow. Only these paths are editable: 'edges.<i>.cond', \
2117'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2118holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2119(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2120— full replacement source for that tool. It is applied only after a recorded \
2121evaluation run passes, so propose one only when the evidence shows the current \
2122code is the defect. \
2123The subject of a fact, the name of a query, the plan hash and the tool name \
2124all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2125becomes a change a human reviewer may apply, so it must be fully supported by \
2126the cited evidence. Propose nothing you cannot ground in the evidence.";
2127
2128const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2134fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2135the facts it relies on are actually present in the cited evidence, NOT that its \
2136conclusion is stated verbatim. Decompose the finding into the factual claims it \
2137depends on. Mark supported=true when those facts are present in the evidence \
2138(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2139on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2140different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2141
2142const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2144each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2145never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2146SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2147'possible' findings with no concrete defect, and reject any claimed \
2148inconsistency or contradiction that is not backed by at least two actually \
2149conflicting facts in the cited evidence. (2) Context — does the finding \
2150correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2151finding when it names a genuine, specific problem grounded in its evidence and \
2152materially useful to a human reviewer; otherwise reject it, and default to \
2153keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2154redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2155grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2156KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2157{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2158
2159const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2161guidance note to help a human reviewer decide. Do not restate the finding. Return \
2162JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2163
2164fn push_evidence(
2170 evidence: &mut Vec<crate::llm::EvidenceItem>,
2171 bundle: &mut BTreeSet<String>,
2172 ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2173 g: &GrainRecord,
2174) {
2175 if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2176 ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2177 evidence.push(crate::llm::EvidenceItem {
2178 hash: g.hash.clone(),
2179 grain_type: g.grain_type.clone(),
2180 text: crate::llm::cap(&grain_brief(g), 400),
2181 });
2182 }
2183}
2184
2185fn grain_brief(g: &GrainRecord) -> String {
2187 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2188 return format!("{s} {r} {o}");
2189 }
2190 if let Some(t) = g.tool_name() {
2195 let status = if g.is_error() { "error" } else { "ok" };
2196 let out = g.tool_content().unwrap_or("");
2197 return format!("tool {t} {status}: {out}");
2198 }
2199 for key in ["content", "body", "text", "summary", "object"] {
2207 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2208 if !v.is_empty() {
2209 return v.to_string();
2210 }
2211 }
2212 }
2213 String::new()
2214}
2215
2216fn sanitize_lesson(s: &str) -> String {
2221 sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2222}
2223
2224fn sanitize_line(s: &str, max: usize) -> String {
2230 let cleaned: String =
2231 s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2232 crate::llm::cap(cleaned.trim(), max)
2233}
2234
2235fn sanitize_relation(s: &str) -> Option<String> {
2239 let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2240 if r.is_empty()
2241 || !r
2242 .chars()
2243 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2244 {
2245 return None;
2246 }
2247 Some(r)
2248}
2249
2250fn safe_definition_body(body: &str) -> bool {
2266 if body.contains('{') || body.contains('}') {
2267 return false;
2268 }
2269 !body
2270 .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2271 .any(|tok| {
2272 ["FORGET", "PURGE", "DROP", "DEFINE"]
2273 .iter()
2274 .any(|kw| tok.eq_ignore_ascii_case(kw))
2275 })
2276}
2277
2278fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2283 let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2284 match resolved {
2285 Some(r) => format!("{summary} {}", r.rendered),
2286 None => summary,
2287 }
2288}
2289
2290struct ValidatedDraft {
2296 draft: crate::llm::LlmDraft,
2297 target_ref: String,
2298 cited: Vec<String>,
2299 resolved: Option<ResolvedProposal>,
2300}
2301
2302struct ResolvedProposal {
2306 action: ActionKind,
2307 proposal: Proposal,
2308 rendered: String,
2311 summary_key: &'static str,
2312 summary_args: serde_json::Map<String, Value>,
2313 rollbackable: bool,
2314 evalset_hash: Option<String>,
2315 importance: f64,
2316 fact_fields: Option<serde_json::Map<String, Value>>,
2320}
2321
2322fn plan_edit_allowed(path: &str) -> bool {
2332 let seg: Vec<&str> = path.split('.').collect();
2333 match seg.as_slice() {
2334 ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2335 ["retries", node] => !node.is_empty(),
2336 _ => false,
2337 }
2338}
2339
2340fn plan_get(body: &Value, path: &str) -> Value {
2343 let mut cur = body;
2344 for seg in path.split('.') {
2345 cur = match cur {
2346 Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2347 Some(v) => v,
2348 None => return Value::Null,
2349 },
2350 Value::Object(o) => match o.get(seg) {
2351 Some(v) => v,
2352 None => return Value::Null,
2353 },
2354 _ => return Value::Null,
2355 };
2356 }
2357 cur.clone()
2358}
2359
2360fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2363 let segs: Vec<&str> = path.split('.').collect();
2364 let Some((last, parents)) = segs.split_last() else {
2365 return false;
2366 };
2367 let mut cur = body;
2368 for seg in parents {
2369 cur = match cur {
2370 Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2371 Some(v) => v,
2372 None => return false,
2373 },
2374 Value::Object(o) => match o.get_mut(*seg) {
2375 Some(v) => v,
2376 None => return false,
2377 },
2378 _ => return false,
2379 };
2380 }
2381 match cur {
2382 Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2383 Some(slot) => {
2384 *slot = to;
2385 true
2386 }
2387 None => false,
2388 },
2389 Value::Object(o) => {
2390 o.insert((*last).to_string(), to);
2391 true
2392 }
2393 _ => false,
2394 }
2395}
2396
2397fn plan_value_ok(path: &str, to: &Value) -> bool {
2402 let seg: Vec<&str> = path.split('.').collect();
2403 match seg.as_slice() {
2404 ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
2405 !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
2406 }),
2407 ["edges", _, "max_cycles"] | ["retries", _] => {
2408 to.as_u64().is_some_and(|n| n <= 1_000)
2409 }
2410 _ => false,
2411 }
2412}
2413
2414fn resolve_proposal<S: OmsSubstrate>(
2420 sub: &S,
2421 d: &crate::llm::LlmDraft,
2422 target: &TargetRef,
2423 cited: &[String],
2424 ns_by_hash: &std::collections::BTreeMap<String, String>,
2425 caps: Capabilities,
2426) -> Option<ResolvedProposal> {
2427 use crate::llm::DraftProposal as P;
2428 let mut args = serde_json::Map::new();
2429 match d.parsed_proposal()? {
2430 P::Lesson { lesson } => {
2432 let lesson = sanitize_lesson(&lesson);
2433 if lesson.is_empty() {
2434 return None;
2435 }
2436 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
2437 args.insert("lesson".into(), Value::from(lesson.clone()));
2438 Some(ResolvedProposal {
2439 action: ActionKind::ClusterFailure,
2443 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2444 rendered: format!("Proposed lesson to record: \"{lesson}\""),
2445 summary_key: "llm.lesson",
2446 summary_args: args,
2447 rollbackable: true,
2448 evalset_hash: None,
2449 importance: 0.5,
2450 fact_fields: Some(fields),
2451 })
2452 }
2453 P::Fact { relation, object } => {
2455 let relation = sanitize_relation(&relation)?;
2456 let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
2457 if object.is_empty() {
2458 return None;
2459 }
2460 let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
2461 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
2462 args.insert("relation".into(), Value::from(relation.clone()));
2463 args.insert("object".into(), Value::from(object.clone()));
2464 Some(ResolvedProposal {
2465 action: ActionKind::Record,
2466 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2467 rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
2468 summary_key: "llm.fact",
2469 summary_args: args,
2470 rollbackable: true,
2471 evalset_hash: None,
2472 importance: 0.5,
2473 fact_fields: Some(fields),
2474 })
2475 }
2476 P::QueryRevision { body } => {
2478 let name = target.opaque();
2479 if name.is_empty()
2483 || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
2484 {
2485 return None;
2486 }
2487 let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
2488 if body.is_empty() || !safe_definition_body(&body) {
2489 return None;
2490 }
2491 let stmt = match target.scheme() {
2492 "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
2493 "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
2494 _ => return None,
2495 };
2496 sub.validate_cal(&stmt).ok()?;
2501 sub.definition_inverse(&stmt).ok().flatten()?;
2502 args.insert("name".into(), Value::from(name));
2503 args.insert("body".into(), Value::from(body.clone()));
2504 Some(ResolvedProposal {
2505 action: ActionKind::Revise,
2506 proposal: Proposal::Cal { cal: stmt },
2507 rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
2508 summary_key: "llm.query_revision",
2509 summary_args: args,
2510 rollbackable: true,
2511 evalset_hash: None,
2512 importance: 0.6,
2513 fact_fields: None,
2514 })
2515 }
2516 P::PlanRevision { edits } => {
2518 if !caps.plans
2519 || target.scheme() != "grain"
2520 || edits.is_empty()
2521 || edits.len() > crate::llm::MAX_PLAN_EDITS
2522 {
2523 return None;
2524 }
2525 let hash = target.opaque();
2526 let g = sub.grain(hash).ok().flatten()?;
2527 if g.grain_type != "workflow" || !g.is_live() {
2528 return None;
2529 }
2530 let mut body = Value::Object(g.fields.clone());
2531 let mut deltas = Vec::new();
2532 let nodes: std::collections::BTreeSet<String> = body
2533 .get("nodes")
2534 .and_then(Value::as_array)
2535 .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
2536 .unwrap_or_default();
2537 for e in &edits {
2538 if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
2539 return None;
2540 }
2541 if let Some(node) = e.path.strip_prefix("retries.") {
2546 if !nodes.contains(node) {
2547 return None;
2548 }
2549 }
2550 if plan_get(&body, &e.path) != e.from {
2553 return None;
2554 }
2555 if e.from == e.to {
2559 return None;
2560 }
2561 if !plan_set(&mut body, &e.path, e.to.clone()) {
2562 return None;
2563 }
2564 deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
2565 }
2566 sub.validate_plan(&body).ok()?;
2570 let Value::Object(fields) = body else {
2571 return None;
2572 };
2573 let stmt = cal::supersede(hash, "workflow", &fields);
2574 sub.validate_cal(&stmt).ok()?;
2579 args.insert("plan".into(), Value::from(hash));
2580 args.insert("edits".into(), Value::from(deltas.join("; ")));
2581 Some(ResolvedProposal {
2582 action: ActionKind::Revise,
2583 proposal: Proposal::Cal { cal: stmt },
2584 rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
2585 summary_key: "llm.plan_revision",
2586 summary_args: args,
2587 rollbackable: true,
2588 evalset_hash: None,
2589 importance: 0.7,
2590 fact_fields: None,
2591 })
2592 }
2593 P::CodeRevision { source } => {
2595 if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
2596 return None;
2597 }
2598 if source.chars().count() > crate::llm::MAX_CODE_LEN {
2599 return None;
2600 }
2601 let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
2605 let mut data = serde_json::Map::new();
2606 data.insert("tool".into(), Value::from(target.opaque()));
2607 data.insert("source".into(), Value::from(source.clone()));
2608 args.insert("tool".into(), Value::from(target.opaque()));
2609 args.insert("bytes".into(), Value::from(source.len() as u64));
2610 Some(ResolvedProposal {
2611 action: ActionKind::CodeRevision,
2612 proposal: Proposal::Data { data },
2613 rendered: format!(
2614 "Proposed new source for tool {} ({} bytes), gated by evalset {}",
2615 target.opaque(),
2616 source.len(),
2617 evalset
2618 ),
2619 summary_key: "llm.code_revision",
2620 summary_args: args,
2621 rollbackable: true,
2622 evalset_hash: Some(evalset),
2623 importance: 0.8,
2624 fact_fields: None,
2625 })
2626 }
2627 }
2628}
2629
2630fn stamp_llm(
2640 model: &str,
2641 d: &crate::llm::LlmDraft,
2642 target_ref: String,
2643 cited: Vec<String>,
2644 resolved: Option<ResolvedProposal>,
2645 confidence: f64,
2646 now_ms: i64,
2647) -> Recommendation {
2648 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2649 let guidance = if d.guidance.trim().is_empty() {
2650 None
2651 } else {
2652 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
2653 };
2654 let (action, proposal, summary, rollbackable, importance, evalset_hash) = match resolved {
2655 Some(mut r) => {
2656 if let Some(mut fields) = r.fact_fields.take() {
2659 fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
2660 r.proposal = Proposal::Cal { cal: cal::add("fact", &fields) };
2661 }
2662 let mut args = r.summary_args;
2663 args.insert("text".into(), Value::from(summary_text));
2664 (
2665 r.action,
2666 r.proposal,
2667 Summary::new(r.summary_key, args),
2668 r.rollbackable,
2669 r.importance,
2670 r.evalset_hash,
2671 )
2672 }
2673 None => {
2674 let mut args = serde_json::Map::new();
2675 args.insert("text".into(), Value::from(summary_text));
2676 let mut data = serde_json::Map::new();
2677 data.insert("source".into(), Value::from("llm"));
2678 (
2679 ActionKind::Flag,
2680 Proposal::Data { data },
2681 Summary::new("llm.discover", args),
2682 false,
2683 0.3,
2684 None,
2685 )
2686 }
2687 };
2688 Recommendation {
2689 hash: String::new(),
2690 analyzer: "loop.llm/1".to_string(),
2691 params_snapshot: serde_json::Map::new(),
2692 origin: Origin::Llm { model: model.to_string() },
2693 target_ref: target_ref.clone(),
2694 action_kind: action,
2695 dedup_key: dedup_key("llm", &target_ref, action),
2696 summary,
2697 severity: Severity::Low,
2698 proposal,
2699 destructive: false,
2700 rollbackable,
2701 evidence: cited,
2702 evidence_query: None,
2703 metric: None,
2704 confidence: confidence.clamp(0.0, 1.0),
2706 importance,
2707 created_at_ms: now_ms,
2708 guidance,
2709 evalset_hash,
2710 status: RecStatus::Pending,
2711 }
2712}
2713
2714fn derived_fact_fields(
2722 target: &TargetRef,
2723 relation: &str,
2724 object: &str,
2725 cited: &[String],
2726 ns_by_hash: &std::collections::BTreeMap<String, String>,
2727) -> Option<serde_json::Map<String, Value>> {
2728 if target.scheme() != "entity" {
2729 return None;
2730 }
2731 let subject = target
2732 .opaque()
2733 .rsplit_once('/')
2734 .map(|(_, s)| s)
2735 .unwrap_or(target.opaque());
2736 if subject.is_empty() {
2737 return None;
2738 }
2739 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
2740 for h in cited {
2741 if let Some(ns) = ns_by_hash.get(h) {
2742 if !ns.is_empty() {
2743 *ns_counts.entry(ns.as_str()).or_default() += 1;
2744 }
2745 }
2746 }
2747 let lesson_ns = ns_counts
2748 .iter()
2749 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
2750 .map(|(ns, _)| ns.to_string());
2751 let mut fields = serde_json::Map::new();
2752 fields.insert("subject".into(), Value::from(subject));
2753 fields.insert("relation".into(), Value::from(relation));
2754 fields.insert("object".into(), Value::from(object));
2755 if let Some(ns) = lesson_ns {
2759 fields.insert("namespace".into(), Value::from(ns));
2760 }
2761 Some(fields)
2762}
2763
2764fn requires_gating(kind: ActionKind) -> bool {
2769 matches!(
2770 kind,
2771 ActionKind::CodeRevision | ActionKind::AdapterRevision
2772 )
2773}
2774
2775fn stamp(
2776 m: &AnalyzerManifest,
2777 params: &crate::manifest::Params,
2778 d: crate::recommendation::RecDraft,
2779 now_ms: i64,
2780) -> Result<Recommendation> {
2781 let target = TargetRef::parse(&d.target_ref)?;
2782 crate::recommendation::validate_code_rules(
2786 d.action_kind,
2787 target.target_class(),
2788 d.evalset_hash.as_deref(),
2789 )?;
2790 let dedup = dedup_key(m.family(), &d.target_ref, d.action_kind);
2791 let destructive = match &d.proposal {
2792 Proposal::Cal { cal } => cal::contains_destructive(cal),
2793 _ => false,
2794 };
2795 let rollbackable = match &d.proposal {
2796 Proposal::Cal { .. } => !destructive,
2797 Proposal::Edit { .. } => false,
2798 Proposal::Data { .. } => requires_gating(d.action_kind),
2802 };
2803 let mut evidence = d.evidence;
2804 evidence.truncate(MAX_EVIDENCE);
2805 let origin = match m.trust_class {
2810 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
2811 _ => Origin::Builtin,
2812 };
2813 Ok(Recommendation {
2814 hash: String::new(),
2815 analyzer: m.id.clone(),
2816 params_snapshot: params.snapshot(),
2817 origin,
2818 target_ref: target.as_string(),
2819 action_kind: d.action_kind,
2820 dedup_key: dedup,
2821 summary: d.summary,
2822 severity: d.severity,
2823 proposal: d.proposal,
2824 destructive,
2825 rollbackable,
2826 evidence,
2827 evidence_query: d.evidence_query,
2828 metric: d.metric,
2829 confidence: d.confidence,
2830 importance: d.importance,
2831 created_at_ms: now_ms,
2832 guidance: None,
2833 evalset_hash: d.evalset_hash,
2834 status: RecStatus::Pending,
2835 })
2836}
2837
2838fn validate_because(because: &str) -> Result<String> {
2839 let trimmed = because.trim();
2840 if trimmed.is_empty() {
2841 return Err(Error::InvalidProposal(
2842 "a BECAUSE reason is required".into(),
2843 ));
2844 }
2845 if trimmed.chars().count() > MAX_BECAUSE {
2846 return Err(Error::InvalidProposal(format!(
2847 "BECAUSE exceeds {MAX_BECAUSE} chars"
2848 )));
2849 }
2850 Ok(trimmed.to_string())
2851}
2852
2853fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
2854 for req in &m.requires {
2855 match req {
2856 Capability::Forks if !caps.forks => return Some("forks"),
2857 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
2858 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
2859 _ => {}
2860 }
2861 }
2862 None
2863}
2864
2865fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
2866 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
2867}
2868
2869fn gate(
2870 opts: &RunOptions,
2871 p: &LoopPersisted,
2872 new_grains: u64,
2873 new_errors: u64,
2874 now_ms: i64,
2875) -> Option<SkipReason> {
2876 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
2877 if !any {
2878 return None;
2879 }
2880 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
2881 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
2882 let stale_ok = opts
2883 .if_stale_ms
2884 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
2885 if min_new_ok || min_err_ok || stale_ok {
2886 return None;
2887 }
2888 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
2890 Some(SkipReason::NotStale)
2891 } else {
2892 Some(SkipReason::MinNewNotMet)
2893 }
2894}
2895
2896fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
2897 let opts = ReadOpts {
2898 live_only: false,
2899 since_ms: watermark.map(|w| w + 1),
2900 };
2901 let mut new_grains = 0u64;
2902 let mut new_errors = 0u64;
2903 for t in [
2904 crate::model::grain_type::FACT,
2905 crate::model::grain_type::EVENT,
2906 crate::model::grain_type::TOOL,
2907 crate::model::grain_type::OBSERVATION,
2908 ] {
2909 let g = sub.grains_of_type(t, None, opts)?;
2910 new_grains += g.len() as u64;
2911 if t == crate::model::grain_type::TOOL {
2913 new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
2914 }
2915 }
2916 Ok((new_grains, new_errors))
2917}
2918
2919fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
2920 let grains = sub.grains_of_type(
2921 crate::model::grain_type::RECOMMENDATION,
2922 Some(LOOP_NS),
2923 ReadOpts {
2924 live_only: false,
2925 since_ms: None,
2926 },
2927 )?;
2928 let mut set = BTreeSet::new();
2929 for g in grains {
2930 let status = p
2931 .status_index
2932 .get(&g.hash)
2933 .copied()
2934 .unwrap_or(RecStatus::Pending);
2935 if matches!(
2940 status,
2941 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
2942 ) {
2943 if let Some(key) = g.str_field("dedup_key") {
2944 set.insert(key.to_string());
2945 }
2946 }
2947 }
2948 Ok(set)
2949}
2950
2951fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
2952 let g = sub
2953 .grain(rec_hash)?
2954 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
2955 Recommendation::from_fields(rec_hash, &g.fields)
2956}
2957
2958pub(crate) fn is_definition_statement(line: &str) -> bool {
2966 let up = line.trim_start().to_ascii_uppercase();
2967 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
2968}
2969
2970const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
2973 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
2974 approve it to acknowledge it and let it expire.";
2975
2976const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
2981 (evalset hash + run id + stats) — use apply_gated";
2982
2983const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
2984 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
2985 acknowledge it and let it expire.";
2986
2987pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
3000 match proposal {
3001 Proposal::Cal { .. } => Ok(()),
3002 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
3003 Proposal::Data { data } => {
3009 if requires_gating(action_kind)
3010 || data.get("revert_of").and_then(Value::as_str).is_some()
3011 {
3012 Ok(())
3013 } else {
3014 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
3015 }
3016 }
3017 }
3018}
3019
3020#[cfg(test)]
3021mod definition_body_tests {
3022 use super::safe_definition_body;
3023
3024 #[test]
3025 fn ordinary_bodies_pass() {
3026 assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
3027 assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
3028 }
3029
3030 #[test]
3031 fn a_body_cannot_close_its_own_block_or_carry_destruction() {
3032 assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
3037 assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
3038 assert!(!safe_definition_body("RECALL facts FORGET abc"));
3040 assert!(!safe_definition_body("recall facts purge older than 1d"));
3041 assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
3042 assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
3044 assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
3046 }
3047}
3048
3049#[cfg(test)]
3050mod plan_edit_tests {
3051 use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
3052 use serde_json::json;
3053
3054 fn plan() -> serde_json::Value {
3055 json!({
3056 "nodes": ["fetch", "review", "post"],
3057 "edges": [
3058 {"src": "fetch", "dst": "review"},
3059 {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
3060 ],
3061 "bindings": {"fetch": "sha256:tool1"},
3062 "retries": {"fetch": 1}
3063 })
3064 }
3065
3066 #[test]
3067 fn the_allowlist_admits_thresholds_and_refuses_topology() {
3068 assert!(plan_edit_allowed("edges.1.cond"));
3069 assert!(plan_edit_allowed("edges.1.max_cycles"));
3070 assert!(plan_edit_allowed("retries.fetch"));
3071 for path in [
3074 "nodes",
3075 "nodes.0",
3076 "edges.0.src",
3077 "edges.0.dst",
3078 "edges",
3079 "bindings.fetch",
3080 "edges.x.cond",
3081 "",
3082 ] {
3083 assert!(!plan_edit_allowed(path), "{path} must not be editable");
3084 }
3085 }
3086
3087 #[test]
3088 fn values_are_type_checked_against_the_field() {
3089 assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
3092 assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
3093 assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
3094 assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
3095 assert!(plan_value_ok("retries.fetch", &json!(3)));
3096 assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
3097 assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
3098 assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
3099 assert!(!plan_value_ok("edges.0.src", &json!("other")));
3100 }
3101
3102 #[test]
3103 fn get_reads_through_arrays_and_objects_and_absence_is_null() {
3104 let p = plan();
3105 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
3106 assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
3107 assert_eq!(plan_get(&p, "retries.review"), json!(null));
3110 assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
3111 assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
3112 }
3113
3114 #[test]
3115 fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
3116 let mut p = plan();
3117 assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
3118 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
3119 assert!(plan_set(&mut p, "retries.review", json!(2)));
3120 assert_eq!(plan_get(&p, "retries.review"), json!(2));
3121 assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
3122 assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
3123 }
3124}
3125
3126#[cfg(test)]
3127mod definition_proposal_tests {
3128 use super::is_definition_statement;
3129
3130 #[test]
3131 fn definition_statements_are_recognized_in_both_spellings() {
3132 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
3133 assert!(is_definition_statement(" define template foo AS { x }"));
3134 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
3135 assert!(!is_definition_statement("ADD fact {}"));
3137 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
3138 assert!(!is_definition_statement("FORGET abc"));
3139 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
3143 }
3144}