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, Checkpoint};
19use crate::substrate::{Capabilities, OmsSubstrate, ReadOpts, SubstrateRead};
20use serde::{Deserialize, Serialize};
21use serde_json::{Map, Value};
22use std::collections::{BTreeMap, BTreeSet};
23
24pub const LOOP_NS: &str = "areev-loop";
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum Scope {
30 Read,
31 Write,
32 Review,
33 Apply,
34 Admin,
35}
36
37#[derive(Debug, Clone, Default)]
39pub struct ScopeSet(Vec<Scope>);
40
41impl ScopeSet {
42 pub fn of(scopes: &[Scope]) -> Self {
43 ScopeSet(scopes.to_vec())
44 }
45 pub fn all() -> Self {
48 ScopeSet(vec![Scope::Admin])
49 }
50 pub fn has(&self, s: Scope) -> bool {
51 self.0.contains(&Scope::Admin) || self.0.contains(&s)
52 }
53}
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub enum Decision {
58 Approve,
59 Reject,
60}
61
62#[derive(Debug, Clone, Default)]
64pub struct RunOptions {
65 pub min_new: Option<u64>,
66 pub min_new_errors: Option<u64>,
67 pub if_stale_ms: Option<i64>,
68 pub namespaces: Vec<String>,
70 pub full_sweep: bool,
77 pub triggering_actor: Option<String>,
84}
85
86#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
88#[serde(rename_all = "lowercase")]
89pub enum RunOutcome {
90 Ran,
91 Skipped,
92}
93
94#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
97#[serde(rename_all = "snake_case")]
98pub enum SkipReason {
99 MinNewNotMet,
100 NotStale,
101 LockHeld,
102 CadenceNotDue,
104}
105
106#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
108pub struct AnalyzerSkip {
109 pub id: String,
110 pub reason: String,
111}
112
113#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
115pub struct RunResult {
116 pub outcome: RunOutcome,
117 #[serde(skip_serializing_if = "Option::is_none")]
118 pub skip_reason: Option<SkipReason>,
119 pub new_grains: u64,
120 pub new_error_events: u64,
121 pub proposed: u64,
122 pub deduped: u64,
123 pub stored: u64,
124 #[serde(default)]
126 pub auto_applied: u64,
127 #[serde(default)]
128 pub analyzers_run: Vec<String>,
129 #[serde(default)]
130 pub analyzers_skipped: Vec<AnalyzerSkip>,
131 #[serde(default, skip_serializing_if = "Option::is_none")]
134 pub llm_funnel: Option<LlmFunnel>,
135 #[serde(default)]
139 pub withdrawn: u64,
140 #[serde(default, skip_serializing_if = "Option::is_none")]
145 pub decider: Option<crate::decide::DeciderReport>,
146}
147
148#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
159pub struct LlmFunnel {
160 pub evidence: u64,
162 pub proposed: u64,
164 pub cited: u64,
166 pub dropped_uncited: u64,
170 pub dropped_target: u64,
173 pub grounded: u64,
175 pub ground_verdicts: u64,
180 pub ground_call_failed: bool,
184 pub kept: u64,
186 pub stored: u64,
188 #[serde(default, skip_serializing_if = "is_zero")]
192 pub advisory_thin_evidence: u64,
193 #[serde(default, skip_serializing_if = "is_zero")]
199 pub dropped_near_duplicate: u64,
200}
201
202fn is_zero(n: &u64) -> bool {
203 *n == 0
204}
205
206impl RunResult {
207 fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
208 RunResult {
209 outcome: RunOutcome::Skipped,
210 skip_reason: Some(reason),
211 new_grains,
212 new_error_events,
213 proposed: 0,
214 deduped: 0,
215 stored: 0,
216 auto_applied: 0,
217 llm_funnel: None,
218 analyzers_run: vec![],
219 analyzers_skipped: vec![],
220 withdrawn: 0,
221 decider: None,
222 }
223 }
224
225 pub fn ran(&self) -> bool {
226 self.outcome == RunOutcome::Ran
227 }
228}
229
230pub struct Engine {
233 analyzers: Vec<Box<dyn Analyzer>>,
234 policy: crate::policy::Policy,
235 llm: Option<Box<dyn crate::llm::LlmBackend>>,
238 ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
243 decider: Option<crate::decide::Decider>,
246}
247
248pub(crate) struct AnalysisPass {
249 pub(crate) survivors: Vec<Recommendation>,
250 proposed: u64,
251 deduped: u64,
252 analyzers_run: Vec<String>,
253 pub(crate) analyzers_skipped: Vec<AnalyzerSkip>,
254 llm_funnel: Option<LlmFunnel>,
255 decider: Option<crate::decide::DeciderReport>,
256}
257
258impl Engine {
259 pub fn with_builtins() -> Self {
262 Engine {
263 analyzers: crate::analyzer::builtin_analyzers(),
264 policy: crate::policy::Policy::default(),
265 llm: None,
266 ground_llm: None,
267 decider: None,
268 }
269 }
270
271 pub fn empty() -> Self {
273 Engine {
274 analyzers: vec![],
275 policy: crate::policy::Policy::default(),
276 llm: None,
277 ground_llm: None,
278 decider: None,
279 }
280 }
281
282 pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
284 self.policy = policy;
285 self
286 }
287
288 pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
293 self.llm = Some(backend);
294 self
295 }
296
297 pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
301 self.ground_llm = Some(backend);
302 self
303 }
304
305 pub fn with_decider(mut self, backend: Box<dyn crate::decide::DecideBackend>) -> Self {
325 let cap = self.decider.as_ref().map(|d| d.pair_cap());
326 let mut d = crate::decide::Decider::new(backend);
327 if let Some(cap) = cap {
328 d = d.with_pair_cap(cap);
329 }
330 self.decider = Some(d);
331 self
332 }
333
334 pub fn with_decider_pair_cap(mut self, cap: usize) -> Self {
339 self.decider = self.decider.take().map(|d| d.with_pair_cap(cap));
340 self
341 }
342
343 pub fn decider(&self) -> Option<&crate::decide::Decider> {
345 self.decider.as_ref()
346 }
347
348 pub fn policy(&self) -> &crate::policy::Policy {
349 &self.policy
350 }
351
352 pub fn has_llm(&self) -> bool {
354 self.llm.is_some()
355 }
356
357 pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
359 self.analyzers.push(analyzer);
360 }
361
362 pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
363 &self.analyzers
364 }
365
366 pub fn analyze_only<S: OmsSubstrate>(
375 &self,
376 sub: &S,
377 opts: &RunOptions,
378 overrides: &BTreeMap<String, Map<String, Value>>,
379 now_ms: i64,
380 ) -> Result<Vec<Recommendation>> {
381 let persisted = LoopPersisted::from_value(sub.load_state()?)?;
382 let analysis_watermark = if opts.full_sweep {
383 None
384 } else {
385 persisted.state.watermark_ms
386 };
387 Ok(self
388 .analysis_pass(
389 sub,
390 &persisted,
391 opts,
392 overrides,
393 analysis_watermark,
394 now_ms,
395 &[],
396 )?
397 .survivors)
398 }
399
400 pub fn run<S: OmsSubstrate>(
403 &self,
404 sub: &mut S,
405 opts: &RunOptions,
406 now_ms: i64,
407 ) -> Result<RunResult> {
408 let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
409 let watermark = persisted.state.watermark_ms;
410 let analysis_watermark = if opts.full_sweep { None } else { watermark };
415
416 let new = count_new(sub, watermark)?;
417 let (new_grains, new_error_events) = (new.grains, new.error_events);
418 if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
419 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
420 }
421 let flags_set = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
425 if !flags_set && !opts.full_sweep {
426 if let Some(reason) = cadence_gate(&self.policy.cadence, &persisted, new, now_ms) {
427 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
428 }
429 }
430
431 let mut outcome_inputs = measure_outcomes(sub, &mut persisted, &self.policy, now_ms)?;
436 let mut withdrawn = 0u64;
437 if self.policy.premise_drift {
438 outcome_inputs.extend(detect_premise_drift(sub, &mut persisted, now_ms)?);
439 withdrawn = withdraw_drifted_open(
442 sub,
443 &mut persisted,
444 self.policy.premise_drift_open_all,
445 now_ms,
446 )?;
447 }
448
449 let AnalysisPass {
450 survivors,
451 proposed,
452 deduped,
453 analyzers_run,
454 analyzers_skipped,
455 llm_funnel,
456 decider,
457 } = self.analysis_pass(
458 &*sub,
459 &persisted,
460 opts,
461 &BTreeMap::new(),
462 analysis_watermark,
463 now_ms,
464 &outcome_inputs,
465 )?;
466
467 let mut stored = 0u64;
470 let mut auto_applied = 0u64;
471 for mut rec in survivors {
472 let spec = rec.to_grain_spec(LOOP_NS)?;
473 let hash = sub.put_grain(&spec)?;
474 rec.hash = hash.clone();
475 let actor = format!("engine:{}", rec.analyzer);
476 let audit = AuditRecord {
477 rec_hash: hash.clone(),
478 from: None,
479 to: RecStatus::Pending,
480 actor: actor.clone(),
481 observer_type: ObserverType::System,
482 because: "analyzer proposed".into(),
483 previous_audit_hash: None,
484 gating: None,
485 at_ms: now_ms,
486 };
487 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
488 persisted
489 .status_index
490 .insert(hash.clone(), RecStatus::Pending);
491 persisted.creators.insert(hash.clone(), actor);
492 if !matches!(rec.origin, Origin::Builtin) {
497 if let Some(trigger) = &opts.triggering_actor {
498 persisted.co_creators.insert(hash.clone(), trigger.clone());
499 }
500 }
501 persisted.audit_heads.insert(hash.clone(), audit_hash);
502 stored += 1;
503
504 if self.can_auto_apply(&*sub, &rec) {
505 self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
506 auto_applied += 1;
507 }
508 }
509
510 persisted.state.last_run_ms = Some(now_ms);
511 persisted.state.watermark_ms = Some(now_ms);
512 sub.store_state(&persisted.to_value()?)?;
513
514 Ok(RunResult {
515 outcome: RunOutcome::Ran,
516 skip_reason: None,
517 new_grains,
518 new_error_events,
519 proposed,
520 deduped,
521 stored,
522 auto_applied,
523 analyzers_run,
524 analyzers_skipped,
525 llm_funnel,
526 withdrawn,
527 decider,
528 })
529 }
530
531 #[allow(clippy::too_many_arguments)]
535 #[allow(clippy::too_many_arguments)]
536 fn analysis_pass<S: OmsSubstrate>(
537 &self,
538 sub: &S,
539 persisted: &LoopPersisted,
540 opts: &RunOptions,
541 external_overrides: &BTreeMap<String, Map<String, Value>>,
542 analysis_watermark: Option<i64>,
543 now_ms: i64,
544 outcome_inputs: &[OutcomeInput],
545 ) -> Result<AnalysisPass> {
546 let existing = existing_dedup_keys(sub, persisted)?;
547 self.analysis_pass_inner(
548 sub,
549 persisted,
550 &self.policy,
551 opts,
552 external_overrides,
553 analysis_watermark,
554 now_ms,
555 outcome_inputs,
556 &existing,
557 None,
558 )
559 }
560
561 #[allow(clippy::too_many_arguments)]
568 pub(crate) fn analysis_pass_inner<S: OmsSubstrate>(
569 &self,
570 sub: &S,
571 persisted: &LoopPersisted,
572 policy: &crate::policy::Policy,
573 opts: &RunOptions,
574 external_overrides: &BTreeMap<String, Map<String, Value>>,
575 analysis_watermark: Option<i64>,
576 now_ms: i64,
577 outcome_inputs: &[OutcomeInput],
578 existing: &BTreeSet<String>,
579 replay: Option<&str>,
580 ) -> Result<AnalysisPass> {
581 let mut analyzers_run = Vec::new();
582 let mut analyzers_skipped = Vec::new();
583 let mut candidates: Vec<Recommendation> = Vec::new();
584 let caps = sub.capabilities();
585 let verdicts = latest_verdicts(persisted);
586 let decider = if replay.is_none() { self.decider.as_ref() } else { None };
590 if let Some(d) = decider {
591 d.reset();
592 }
593
594 for analyzer in &self.analyzers {
595 let m = analyzer.manifest();
596 if let Some(why) = replay {
597 let out_of_process = m.trust_class == crate::manifest::TrustClass::Command;
598 let telemetry_fed = m.requires.contains(&crate::manifest::Capability::Telemetry);
599 if out_of_process || telemetry_fed {
600 analyzers_skipped.push(AnalyzerSkip {
601 id: m.id.clone(),
602 reason: format!(
603 "not replayed: {}",
604 if out_of_process { why } else { "telemetry rollups are not time-indexed" }
605 ),
606 });
607 continue;
608 }
609 }
610 let cfg = persisted.config.get(&m.id);
611 let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
612 if !enabled {
613 analyzers_skipped.push(AnalyzerSkip {
614 id: m.id.clone(),
615 reason: "disabled".into(),
616 });
617 continue;
618 }
619 if policy.denies(m.family()) {
620 analyzers_skipped.push(AnalyzerSkip {
621 id: m.id.clone(),
622 reason: "denied by host policy".into(),
623 });
624 continue;
625 }
626 if let Some(missing) = missing_capability(m, caps) {
627 analyzers_skipped.push(AnalyzerSkip {
628 id: m.id.clone(),
629 reason: format!("missing capability: {missing}"),
630 });
631 continue;
632 }
633 let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
634 if let Some(extra) = external_overrides.get(&m.id) {
635 for (key, value) in extra {
636 param_overrides.insert(key.clone(), value.clone());
637 }
638 }
639 let params = match m.resolve_params(¶m_overrides) {
640 Ok(p) => p,
641 Err(e) => {
642 analyzers_skipped.push(AnalyzerSkip {
643 id: m.id.clone(),
644 reason: e.to_string(),
645 });
646 continue;
647 }
648 };
649 let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
650 let ns_slice: &[String] = if ns_owned.is_empty() {
651 &opts.namespaces
652 } else {
653 &ns_owned
654 };
655 let reader: &dyn SubstrateRead = sub;
656 let ctx = AnalyzeCtx::new(
657 reader,
658 ¶ms,
659 ns_slice,
660 analysis_watermark,
661 now_ms,
662 outcome_inputs,
663 &verdicts,
664 )
665 .with_decider(decider);
666 match analyzer.analyze(&ctx) {
667 Ok(drafts) => {
668 analyzers_run.push(m.id.clone());
669 for draft in drafts {
670 match stamp(m, ¶ms, draft, now_ms, ns_slice) {
671 Ok(rec) => candidates.push(rec),
672 Err(e) => analyzers_skipped.push(AnalyzerSkip {
673 id: m.id.clone(),
674 reason: e.to_string(),
675 }),
676 }
677 }
678 }
679 Err(e) => analyzers_skipped.push(AnalyzerSkip {
680 id: m.id.clone(),
681 reason: e.to_string(),
682 }),
683 }
684 }
685
686 let mut funnel = LlmFunnel::default();
687 if self.llm.is_some() && replay.is_none() {
688 candidates.extend(self.discover(
689 sub,
690 &candidates,
691 analysis_watermark,
692 &opts.namespaces,
693 now_ms,
694 &mut funnel,
695 ));
696 }
697
698 let proposed = candidates.len() as u64;
699 let mut seen = BTreeSet::new();
700 let mut survivors = Vec::new();
701 for candidate in candidates {
702 let family = crate::manifest::analyzer_family(&candidate.analyzer);
703 let floor = [
704 severity_floor_for(persisted, &candidate.analyzer),
705 policy.severity_floor(family),
706 ]
707 .into_iter()
708 .flatten()
709 .max();
710 if floor.is_some_and(|floor| candidate.severity < floor) {
711 continue;
712 }
713 if !seen.insert(candidate.dedup_key.clone()) {
714 continue;
715 }
716 if existing.contains(&candidate.dedup_key) {
717 continue;
718 }
719 if persisted
720 .cooldowns
721 .get(&candidate.dedup_key)
722 .is_some_and(|until| now_ms < *until)
723 {
724 continue;
725 }
726 survivors.push(candidate);
727 }
728 let deduped = proposed - survivors.len() as u64;
729 if self.llm.is_some() && replay.is_none() {
730 self.enrich(&mut survivors);
731 }
732 Ok(AnalysisPass {
733 survivors,
734 proposed,
735 deduped,
736 analyzers_run,
737 analyzers_skipped,
738 llm_funnel: self.llm.is_some().then_some(funnel),
739 decider: decider.map(|d| d.report()),
740 })
741 }
742
743 fn discover<S: OmsSubstrate>(
750 &self,
751 sub: &S,
752 candidates: &[Recommendation],
753 watermark: Option<i64>,
754 namespaces: &[String],
755 now_ms: i64,
756 funnel: &mut LlmFunnel,
757 ) -> Vec<Recommendation> {
758 let Some(llm) = &self.llm else {
759 return Vec::new();
760 };
761 let findings: Vec<crate::llm::FindingBrief> = candidates
762 .iter()
763 .take(32)
764 .map(|c| crate::llm::FindingBrief {
765 analyzer: c.analyzer.clone(),
766 summary: c.summary.render(),
767 target: c.target_ref.clone(),
768 severity: c.severity.as_str().to_string(),
769 })
770 .collect();
771 let attribution = self.policy.evidence_attribution;
777 let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
778 let mut bundle: BTreeSet<String> = BTreeSet::new();
779 let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
780 'cited: for c in candidates {
781 for h in &c.evidence {
782 if evidence.len() >= CITED_SEED_CAP {
791 break 'cited;
792 }
793 if !bundle.contains(h) {
794 if let Ok(Some(g)) = sub.grain(h) {
795 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
796 }
797 }
798 }
799 }
800 let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
801 vec![None]
802 } else {
803 namespaces.iter().map(|n| Some(n.as_str())).collect()
804 };
805 let opts = ReadOpts { live_only: true, since_ms: watermark };
806 let mut tool_seeded = 0usize;
827 let seed_successes = self.policy.skills.enabled || self.policy.plans.enabled;
828 'tools: for want_error in [true, false] {
829 if !want_error && !seed_successes {
830 break;
831 }
832 for ns in &scan_ns {
833 if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
834 for g in recent {
835 if tool_seeded >= TOOL_SEED_CAP
836 || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
837 {
838 break 'tools;
839 }
840 if g.is_error() != want_error {
841 continue;
842 }
843 let before = evidence.len();
844 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
845 if evidence.len() > before {
846 tool_seeded += 1;
847 }
848 }
849 }
850 }
851 }
852 if let Ok(rows) =
861 sub.grains_of_type(crate::model::grain_type::OBSERVATION, Some(crate::eval::HARNESS_NS), opts)
862 {
863 let mut seeded = 0usize;
864 for g in rows {
865 if seeded >= HARNESS_SEED_CAP || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE {
866 break;
867 }
868 let kind = g.str_field("observation_kind").unwrap_or_default();
869 if !HARNESS_EVIDENCE_KINDS.contains(&kind) {
870 continue;
871 }
872 let before = evidence.len();
873 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
874 if evidence.len() > before {
875 seeded += 1;
876 }
877 }
878 }
879 'notes: for ns in &scan_ns {
889 if let Ok(recent) =
890 sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
891 {
892 for g in recent {
893 if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
894 break 'notes;
895 }
896 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
897 }
898 }
899 }
900 'seed: for gt in [
901 crate::model::grain_type::FACT,
902 crate::model::grain_type::OBSERVATION,
903 ] {
904 for ns in &scan_ns {
905 if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
906 for g in recent {
907 if evidence.len() >= EVIDENCE_CAP {
908 break 'seed;
909 }
910 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
911 }
912 }
913 }
914 }
915 funnel.evidence = evidence.len() as u64;
916 if evidence.is_empty() {
917 return Vec::new(); }
919 let (approved, rejected) = self.llm_history(sub);
924 let base = match self.policy.discover_objective {
925 crate::policy::DiscoverObjective::ReviewQueue => DISCOVER_INSTRUCTIONS,
926 crate::policy::DiscoverObjective::Learner => DISCOVER_LEARNER_INSTRUCTIONS,
927 };
928 let mut instructions = base.to_string();
931 if self.policy.skills.enabled {
932 instructions.push_str(&skill_instructions(self.policy.skills.min_steps));
933 }
934 if self.policy.plans.enabled {
935 instructions.push_str(&plan_instructions(self.policy.plans.min_nodes));
936 }
937 if findings.iter().any(|f| f.analyzer.starts_with("loop.lesson_pile/")) {
940 instructions.push_str(CONSOLIDATION_INSTRUCTIONS);
941 }
942 let request = crate::llm::LlmRequest {
943 loop_proto: 1,
944 op: "discover",
945 instructions: &instructions,
946 findings: findings.clone(),
947 evidence: evidence.clone(),
948 rejected,
949 approved,
950 };
951 let Ok(body) = serde_json::to_string(&request) else {
952 return Vec::new();
953 };
954 let raw = match llm.complete(&body) {
955 Ok(r) => r,
956 Err(_) => return Vec::new(), };
958 let caps = sub.capabilities();
962 let mut validated: Vec<ValidatedDraft> = Vec::new();
963 let drafts: Vec<_> = crate::llm::parse_discover(&raw)
964 .recommendations
965 .into_iter()
966 .take(crate::llm::MAX_LLM_DRAFTS)
967 .collect();
968 funnel.proposed = drafts.len() as u64;
969 let id_to_hash: std::collections::BTreeMap<&str, &str> = evidence
970 .iter()
971 .map(|e| (e.id.as_str(), e.hash.as_str()))
972 .collect();
973 for d in drafts {
974 let mut cited: Vec<String> = Vec::new();
975 for c in &d.evidence {
976 if let Some(h) = resolve_citation(c, &bundle, &id_to_hash) {
977 if !cited.contains(&h) {
978 cited.push(h);
979 }
980 }
981 }
982 if cited.is_empty() {
983 funnel.dropped_uncited += 1;
984 continue; }
986 let Ok(target) = TargetRef::parse(&d.target) else {
987 funnel.dropped_target += 1;
988 continue;
989 };
990 let tc = target.target_class();
991 if !matches!(tc, "memory" | "query" | "code") {
997 funnel.dropped_target += 1;
998 continue;
999 }
1000 let thin = cited.len() < self.policy.min_evidence as usize;
1007 if thin {
1008 funnel.advisory_thin_evidence += 1;
1009 }
1010 let resolved = if thin {
1011 None
1012 } else {
1013 resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps, &self.policy)
1014 };
1015 if tc == "code" && resolved.is_none() {
1020 funnel.dropped_target += 1;
1021 continue;
1022 }
1023 validated.push(ValidatedDraft {
1024 draft: d,
1025 target_ref: target.as_string(),
1026 cited,
1027 resolved,
1028 });
1029 }
1030 funnel.cited = validated.len() as u64;
1031 if validated.is_empty() {
1032 return Vec::new();
1033 }
1034 let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
1041 let outcome_metric = self.outcome_metric_template(sub);
1042 self.verify_drafts(
1043 sub, &**llm, ground, validated, &evidence, outcome_metric, now_ms, funnel, namespaces,
1044 )
1045 }
1046
1047 fn outcome_metric_template<S: OmsSubstrate>(
1055 &self,
1056 sub: &S,
1057 ) -> Option<crate::recommendation::MetricSnapshot> {
1058 let e = self.policy.outcome_evalset.as_ref()?;
1059 let run = crate::eval::newest_eval_run(sub, &e.hash, None).ok().flatten()?;
1060 let baseline = crate::eval::run_value(&run, &e.field)?;
1061 let schedule = e.schedule();
1065 let ms_only: Vec<i64> = schedule.iter().filter_map(Checkpoint::as_ms).collect();
1066 let all_ms = ms_only.len() == schedule.len();
1067 let horizons = if ms_only.is_empty() { vec![86_400_000] } else { ms_only };
1068 Some(crate::recommendation::MetricSnapshot {
1069 metric: format!("evalset:{}:{}", e.hash, e.field),
1070 baseline,
1071 unit: e.field.clone(),
1072 n: run.total(),
1073 window: "per-run".into(),
1074 subject: None,
1075 namespace: None,
1076 relation: None,
1077 query: format!(
1078 "RECALL facts WHERE subject = \"evalset:{}\" AND relation = \"mg:eval_run\"",
1079 e.hash
1080 ),
1081 review_after_ms: horizons[0],
1082 horizons_ms: if all_ms { horizons } else { Vec::new() },
1083 checkpoints: if all_ms { Vec::new() } else { schedule },
1084 higher_is_better: e.higher_is_better,
1085 })
1086 }
1087
1088 #[allow(clippy::too_many_arguments)]
1095 #[allow(clippy::too_many_arguments)]
1096 fn verify_drafts<S: SubstrateRead>(
1097 &self,
1098 sub: &S,
1099 llm: &dyn crate::llm::LlmBackend,
1100 ground: &dyn crate::llm::LlmBackend,
1101 validated: Vec<ValidatedDraft>,
1102 evidence: &[crate::llm::EvidenceItem],
1103 outcome_metric: Option<crate::recommendation::MetricSnapshot>,
1104 now_ms: i64,
1105 funnel: &mut LlmFunnel,
1106 scope: &[String],
1109 ) -> Vec<Recommendation> {
1110 use crate::llm::*;
1111 let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
1112 evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
1113 let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
1114 cited
1115 .iter()
1116 .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
1117 .collect()
1118 };
1119
1120 let claims: Vec<GroundItem> = validated
1124 .iter()
1125 .enumerate()
1126 .map(|(i, v)| GroundItem {
1127 id: i,
1128 claim: claim_text(&v.draft, v.resolved.as_ref()),
1129 evidence: ev_for(&v.cited),
1130 })
1131 .collect();
1132 let decided = self.decide_ground_verify(&validated, &claims);
1139 let ground_req = GroundRequest {
1140 loop_proto: 1,
1141 op: "ground",
1142 instructions: GROUND_INSTRUCTIONS,
1143 claims,
1144 };
1145 let grounded: std::collections::BTreeSet<usize> = if let Some(d) = &decided {
1151 funnel.ground_verdicts = d.len() as u64;
1152 d.iter().filter(|(_, j)| j.grounded).map(|(i, _)| *i).collect()
1153 } else {
1154 match serde_json::to_string(&ground_req)
1155 .ok()
1156 .and_then(|b| ground.complete(&b).ok())
1157 {
1158 Some(raw) => {
1159 let parsed = parse_ground(&raw);
1160 funnel.ground_verdicts = parsed.results.len() as u64;
1161 parsed
1162 .results
1163 .into_iter()
1164 .filter(|r| r.supported)
1165 .map(|r| r.id)
1166 .collect()
1167 }
1168 None => {
1169 funnel.ground_call_failed = true;
1170 return Vec::new();
1171 }
1172 }
1173 };
1174 funnel.grounded = grounded.len() as u64;
1175 if grounded.is_empty() {
1176 return Vec::new();
1177 }
1178
1179 let items: Vec<VerifyItem> = validated
1185 .iter()
1186 .enumerate()
1187 .filter(|(i, _)| grounded.contains(i))
1188 .map(|(i, v)| VerifyItem {
1189 id: i,
1190 summary: claim_text(&v.draft, v.resolved.as_ref()),
1192 target: v.target_ref.clone(),
1193 evidence: ev_for(&v.cited),
1194 })
1195 .collect();
1196 let verify_req = VerifyRequest {
1197 loop_proto: 1,
1198 op: "verify",
1199 instructions: VERIFY_INSTRUCTIONS,
1200 findings: items,
1201 };
1202 let verdicts: std::collections::BTreeMap<usize, f64> =
1203 match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
1204 Some(raw) => parse_verify(&raw)
1205 .results
1206 .into_iter()
1207 .filter(|r| r.keep)
1208 .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
1209 .collect(),
1210 None => return Vec::new(),
1211 };
1212
1213 funnel.kept = verdicts.len() as u64;
1214 let mut out = Vec::new();
1218 for (i, v) in validated.into_iter().enumerate() {
1219 if let Some(&self_report) = verdicts.get(&i) {
1220 let judged = decided.as_ref().and_then(|d| d.get(&i));
1224 let conf = judged.map_or(self_report, |j| j.sound);
1225 if conf >= MIN_LLM_CONFIDENCE {
1226 let near = match v.resolved.as_ref() {
1235 Some(r) if r.action != ActionKind::Consolidate => r
1236 .fact_fields
1237 .as_ref()
1238 .filter(|f| f.get("relation").and_then(Value::as_str) == Some("lesson"))
1239 .map(|f| {
1240 near_duplicates_of(
1241 sub,
1242 f.get("subject").and_then(Value::as_str).unwrap_or(""),
1243 f.get("namespace").and_then(Value::as_str),
1244 f.get("object").and_then(Value::as_str).unwrap_or(""),
1245 )
1246 })
1247 .unwrap_or_default(),
1248 _ => Vec::new(),
1249 };
1250 if !near.is_empty()
1251 && self.policy.near_duplicate == crate::policy::NearDuplicateMode::Suppress
1252 {
1253 funnel.dropped_near_duplicate += 1;
1254 continue;
1255 }
1256 let mut rec = stamp_llm(
1257 llm.model(),
1258 &v.draft,
1259 v.target_ref,
1260 v.cited,
1261 v.resolved,
1262 conf,
1263 now_ms,
1264 scope,
1265 );
1266 if let Some(j) = judged {
1267 rec.llm_confidence = Some(self_report);
1270 rec.judged_by = Some(j.judged_by.clone());
1271 }
1272 if let Some(best) = near.first() {
1273 rec.summary.args.insert("near_count".into(), Value::from(near.len() as u64));
1275 rec.summary.args.insert("near_score".into(), Value::from(best.score));
1276 rec.summary.args.insert("near_method".into(), Value::from(best.method.clone()));
1277 rec.summary.args.insert(
1278 "near_hash".into(),
1279 Value::from(best.hash.chars().take(12).collect::<String>()),
1280 );
1281 rec.summary.template_id = "llm.lesson_near_duplicate".into();
1282 rec.near_duplicate_of = near;
1283 }
1284 if rec.rollbackable {
1288 rec.metric = outcome_metric.clone();
1289 }
1290 out.push(rec);
1291 }
1292 }
1293 }
1294 funnel.stored = out.len() as u64;
1295 out
1296 }
1297
1298 fn decide_ground_verify(
1309 &self,
1310 validated: &[ValidatedDraft],
1311 claims: &[crate::llm::GroundItem],
1312 ) -> Option<BTreeMap<usize, DecidedDraft>> {
1313 use crate::decide::{Ask, DECIDE_MIN_P};
1314 let d = self.decider.as_ref()?;
1315 if !d.calibrated() {
1316 return None;
1317 }
1318 let backend = d.describe();
1319 let mut out = BTreeMap::new();
1320 for (i, (v, c)) in validated.iter().zip(claims).enumerate() {
1321 let evidence: Vec<Value> = c
1322 .evidence
1323 .iter()
1324 .map(|e| serde_json::json!({"id": e.id, "grain_type": e.grain_type, "text": e.text}))
1325 .collect();
1326 let state = serde_json::json!({
1327 "recommendation": {"summary": c.claim, "guidance": v.draft.guidance},
1328 "evidence": evidence,
1329 });
1330 let mut asks: Vec<Ask> = c
1331 .evidence
1332 .iter()
1333 .map(|e| Ask::Noul {
1334 id: format!("ev_{}", e.id),
1335 instructions: format!(
1336 "Does evidence item \"{}\" (in state.evidence) contain the premise the recommendation relies on?",
1337 e.id
1338 ),
1339 })
1340 .collect();
1341 asks.push(Ask::Noul {
1342 id: "sound".into(),
1343 instructions: "Given only this evidence, is the recommendation sound?".into(),
1344 });
1345 let a = d.ask(state, &asks).ok().filter(|a| a.calibrated)?;
1346 let sound = *a.noul.get("sound")?;
1347 let grounded = a
1348 .noul
1349 .iter()
1350 .any(|(id, p)| id.starts_with("ev_") && *p >= DECIDE_MIN_P);
1351 let judged_by = a.judged_by(&backend, "ground_verify", a.noul.clone());
1352 out.insert(i, DecidedDraft { grounded, sound, judged_by });
1353 }
1354 Some(out)
1355 }
1356
1357 fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
1362 const MAX: usize = 20;
1363 let Ok(mut recs) = self.recommendations(sub, None) else {
1364 return (Vec::new(), Vec::new());
1365 };
1366 recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
1367 recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
1368 let mut approved = Vec::new();
1369 let mut rejected = Vec::new();
1370 for r in &recs {
1371 match r.status {
1372 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
1373 if approved.len() < MAX =>
1374 {
1375 approved.push(r.summary.render());
1376 }
1377 RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
1378 _ => {}
1379 }
1380 }
1381 (approved, rejected)
1382 }
1383
1384 fn enrich(&self, survivors: &mut [Recommendation]) {
1389 let Some(llm) = &self.llm else {
1390 return;
1391 };
1392 if survivors.is_empty() {
1393 return;
1394 }
1395 let findings: Vec<crate::llm::FindingBrief> = survivors
1396 .iter()
1397 .map(|r| crate::llm::FindingBrief {
1398 analyzer: r.analyzer.clone(),
1399 summary: r.summary.render(),
1400 target: r.target_ref.clone(),
1401 severity: r.severity.as_str().to_string(),
1402 })
1403 .collect();
1404 let request = crate::llm::LlmRequest {
1405 loop_proto: 1,
1406 op: "enrich",
1407 instructions: ENRICH_INSTRUCTIONS,
1408 findings,
1409 evidence: Vec::new(),
1410 rejected: Vec::new(),
1411 approved: Vec::new(),
1412 };
1413 let Ok(body) = serde_json::to_string(&request) else {
1414 return;
1415 };
1416 let raw = match llm.complete(&body) {
1417 Ok(r) => r,
1418 Err(_) => return,
1419 };
1420 for note in crate::llm::parse_enrich(&raw).notes {
1421 if note.guidance.trim().is_empty() {
1422 continue;
1423 }
1424 if let Some(r) = survivors
1425 .iter_mut()
1426 .find(|r| r.target_ref == note.target && r.guidance.is_none())
1427 {
1428 r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
1429 }
1430 }
1431 }
1432
1433 fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
1442 if !rec.origin.auto_apply_eligible() || rec.destructive {
1443 return false;
1444 }
1445 if rec.judged_by.is_some() {
1450 return false;
1451 }
1452 let manifest_ok = self
1456 .analyzers
1457 .iter()
1458 .map(|a| a.manifest())
1459 .find(|m| m.id == rec.analyzer)
1460 .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
1461 if !manifest_ok {
1462 return false;
1463 }
1464 let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1465 return false;
1466 };
1467 let family = crate::manifest::analyzer_family(&rec.analyzer);
1468 if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1469 return false;
1470 }
1471 match &rec.proposal {
1476 Proposal::Cal { cal } => cal
1477 .lines()
1478 .map(str::trim)
1479 .filter(|l| !l.is_empty())
1480 .all(|l| supersede_is_value_identical(sub, l)),
1481 _ => false,
1482 }
1483 }
1484
1485 fn auto_apply<S: OmsSubstrate>(
1488 &self,
1489 sub: &mut S,
1490 p: &mut LoopPersisted,
1491 rec: &Recommendation,
1492 now_ms: i64,
1493 ) -> Result<()> {
1494 let mut created = Vec::new();
1495 if let Proposal::Cal { cal } = &rec.proposal {
1496 if cal.lines().map(str::trim).any(is_definition_statement) {
1501 return Err(Error::InvalidProposal(
1502 "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1503 auto-applied: it changes what every future context contains, so it \
1504 requires a human APPROVE + APPLY with BECAUSE"
1505 .into(),
1506 ));
1507 }
1508 for r in sub.execute_cal(cal)? {
1509 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1510 created.push(h.to_string());
1511 }
1512 }
1513 }
1514 let applied = AppliedRecord {
1515 applied_at_ms: now_ms,
1516 target_ref: rec.target_ref.clone(),
1517 rollbackable: rec.rollbackable,
1518 created_hashes: created,
1519 inverse_cal: None,
1520 metric: rec.metric.clone(),
1521 };
1522 let prev = p.audit_heads.get(&rec.hash).cloned();
1523 let audit = AuditRecord {
1524 rec_hash: rec.hash.clone(),
1525 from: Some(RecStatus::Pending),
1526 to: RecStatus::Applied,
1527 actor: "policy:auto".into(),
1528 observer_type: ObserverType::Policy,
1529 because: "auto-applied per host policy".into(),
1530 previous_audit_hash: prev,
1531 gating: None,
1532 at_ms: now_ms,
1533 };
1534 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1535 p.audit_heads.insert(rec.hash.clone(), audit_hash);
1536 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1537 p.applied.insert(rec.hash.clone(), applied);
1538 Ok(())
1539 }
1540
1541 #[allow(clippy::too_many_arguments)]
1544 pub fn review<S: OmsSubstrate>(
1545 &self,
1546 sub: &mut S,
1547 rec_hash: &str,
1548 decision: Decision,
1549 actor: &str,
1550 observer: ObserverType,
1551 scopes: &ScopeSet,
1552 because: &str,
1553 now_ms: i64,
1554 ) -> Result<()> {
1555 if !scopes.has(Scope::Review) {
1556 return Err(Error::ScopeDenied("review".into()));
1557 }
1558 let because = validate_because(because)?;
1559 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1560 let status = *p
1561 .status_index
1562 .get(rec_hash)
1563 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1564 let to = match decision {
1565 Decision::Approve => RecStatus::Approved,
1566 Decision::Reject => RecStatus::Rejected,
1567 };
1568 if !status.can_transition_to(to, false) {
1569 return Err(Error::LifecycleViolation(format!(
1570 "{} -> {}",
1571 status.as_str(),
1572 to.as_str()
1573 )));
1574 }
1575 if to == RecStatus::Approved {
1576 if let Some(creator) = p.creators.get(rec_hash) {
1577 if creator == actor {
1578 return Err(Error::SelfApproval(format!(
1579 "{actor} created this recommendation"
1580 )));
1581 }
1582 }
1583 if let Some(trigger) = p.co_creators.get(rec_hash) {
1584 if trigger == actor {
1585 return Err(Error::SelfApproval(format!(
1586 "{actor} triggered the run that authored this recommendation"
1587 )));
1588 }
1589 }
1590 }
1591 let prev = p.audit_heads.get(rec_hash).cloned();
1592 let audit = AuditRecord {
1593 rec_hash: rec_hash.into(),
1594 from: Some(status),
1595 to,
1596 actor: actor.into(),
1597 observer_type: observer,
1598 because,
1599 previous_audit_hash: prev,
1600 gating: None,
1601 at_ms: now_ms,
1602 };
1603 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1604 p.audit_heads.insert(rec_hash.into(), audit_hash);
1605 p.status_index.insert(rec_hash.into(), to);
1606 if to == RecStatus::Rejected {
1607 if let Ok(rec) = load_rec(sub, rec_hash) {
1608 strike_cooldown(&mut p, rec.dedup_key, now_ms);
1609 }
1610 }
1611 sub.store_state(&p.to_value()?)?;
1612 Ok(())
1613 }
1614
1615 pub fn preflight_apply<S: OmsSubstrate>(
1632 &self,
1633 sub: &S,
1634 rec_hash: &str,
1635 scopes: &ScopeSet,
1636 allow_destructive: bool,
1637 has_gating: bool,
1638 ) -> Result<()> {
1639 if !scopes.has(Scope::Apply) {
1640 return Err(Error::ScopeDenied("apply".into()));
1641 }
1642 let rec = load_rec(sub, rec_hash)?;
1643 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1644 return Err(Error::DestructiveGated(
1645 "destructive apply requires admin scope + allow_destructive".into(),
1646 ));
1647 }
1648 ensure_executable(rec.action_kind, &rec.proposal)?;
1649 if requires_gating(rec.action_kind) && !has_gating {
1650 return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1651 }
1652 Ok(())
1653 }
1654
1655 #[allow(clippy::too_many_arguments)]
1659 pub fn apply<S: OmsSubstrate>(
1660 &self,
1661 sub: &mut S,
1662 rec_hash: &str,
1663 actor: &str,
1664 observer: ObserverType,
1665 scopes: &ScopeSet,
1666 because: &str,
1667 allow_destructive: bool,
1668 now_ms: i64,
1669 ) -> Result<AppliedRecord> {
1670 self.apply_inner(
1671 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1672 )
1673 }
1674
1675 pub fn gating_evidence<S: OmsSubstrate>(
1682 &self,
1683 sub: &S,
1684 rec_hash: &str,
1685 run_id: &str,
1686 ) -> Result<crate::recommendation::GatingEvidence> {
1687 let rec = self
1688 .recommendations(sub, None)?
1689 .into_iter()
1690 .find(|r| r.hash == rec_hash)
1691 .ok_or_else(|| {
1692 Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1693 })?;
1694 let pin = rec.evalset_hash.ok_or_else(|| {
1695 Error::InvalidProposal(
1696 "this recommendation pins no evalset — a gating run applies only \
1697 to code and adapter revisions"
1698 .into(),
1699 )
1700 })?;
1701 match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1706 Some(run) => Ok(crate::recommendation::GatingEvidence {
1707 evalset_hash: pin,
1708 run_id: run.run_id,
1709 passed: run.passed,
1710 failed: run.failed,
1711 }),
1712 None => Err(Error::InvalidProposal(format!(
1713 "no recorded gate run '{run_id}' for evalset {pin} — run \
1714 `areev eval run --evalset {pin} ...` first"
1715 ))),
1716 }
1717 }
1718
1719 #[allow(clippy::too_many_arguments)]
1723 pub fn apply_gated<S: OmsSubstrate>(
1724 &self,
1725 sub: &mut S,
1726 rec_hash: &str,
1727 actor: &str,
1728 observer: ObserverType,
1729 scopes: &ScopeSet,
1730 because: &str,
1731 allow_destructive: bool,
1732 gating: &crate::recommendation::GatingEvidence,
1733 now_ms: i64,
1734 ) -> Result<AppliedRecord> {
1735 self.apply_inner(
1736 sub,
1737 rec_hash,
1738 actor,
1739 observer,
1740 scopes,
1741 because,
1742 allow_destructive,
1743 Some(gating),
1744 now_ms,
1745 )
1746 }
1747
1748 #[allow(clippy::too_many_arguments)]
1749 fn apply_inner<S: OmsSubstrate>(
1750 &self,
1751 sub: &mut S,
1752 rec_hash: &str,
1753 actor: &str,
1754 observer: ObserverType,
1755 scopes: &ScopeSet,
1756 because: &str,
1757 allow_destructive: bool,
1758 gating: Option<&crate::recommendation::GatingEvidence>,
1759 now_ms: i64,
1760 ) -> Result<AppliedRecord> {
1761 if !scopes.has(Scope::Apply) {
1762 return Err(Error::ScopeDenied("apply".into()));
1763 }
1764 let because = validate_because(because)?;
1765 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1766 let status = *p
1767 .status_index
1768 .get(rec_hash)
1769 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1770 if !status.can_transition_to(RecStatus::Applied, false) {
1771 return Err(Error::LifecycleViolation(format!(
1772 "{} -> applied (approve first)",
1773 status.as_str()
1774 )));
1775 }
1776 let rec = load_rec(sub, rec_hash)?;
1777 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1778 return Err(Error::DestructiveGated(
1779 "destructive apply requires admin scope + allow_destructive".into(),
1780 ));
1781 }
1782 if requires_gating(rec.action_kind) {
1787 let g = gating
1788 .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1789 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1790 if g.evalset_hash != pin {
1791 return Err(Error::InvalidProposal(format!(
1792 "gating ran evalset {} but the recommendation is pinned \
1793 to {pin} (Rule E1)",
1794 g.evalset_hash
1795 )));
1796 }
1797 match sub.grain(pin)? {
1798 Some(evalset) if evalset.is_live() => {}
1799 Some(_) => {
1800 return Err(Error::InvalidProposal(
1801 "the pinned evalset was superseded after gating — \
1802 the recommendation must re-gate (Rule E1)"
1803 .into(),
1804 ))
1805 }
1806 None => {
1807 return Err(Error::InvalidProposal(format!(
1808 "pinned evalset {pin} not found in the substrate"
1809 )))
1810 }
1811 }
1812 if g.failed > 0 {
1813 return Err(Error::InvalidProposal(format!(
1814 "the gating run failed {}/{} cases — a failing gate \
1815 admits nothing",
1816 g.failed,
1817 g.passed + g.failed
1818 )));
1819 }
1820 }
1821
1822 let mut created = Vec::new();
1824 let mut inverse_cal: Option<String> = None;
1828 match &rec.proposal {
1829 Proposal::Cal { cal } => {
1830 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1831 if !is_definition_statement(line) {
1832 continue;
1833 }
1834 match sub.definition_inverse(line)? {
1835 Some(inv) => inverse_cal = Some(inv),
1836 None => {
1837 return Err(Error::InvalidProposal(format!(
1838 "this substrate cannot record a rollback inverse for {line:?}; \
1839 a definition rewrite that ROLLBACK could not undo is refused \
1840 rather than applied"
1841 )))
1842 }
1843 }
1844 }
1845 let rows = sub.execute_cal(cal)?;
1846 for r in rows {
1847 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1848 created.push(h.to_string());
1849 }
1850 }
1851 }
1852 Proposal::Edit { .. } => {
1855 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1856 }
1857 Proposal::Data { data } if requires_gating(rec.action_kind) => {
1865 let relation = if rec.action_kind == ActionKind::AdapterRevision {
1866 "mg:adapter_promotion"
1867 } else {
1868 "mg:code_promotion"
1869 };
1870 let mut promoted = data.clone();
1878 if let Some(Value::String(src)) = promoted.remove("source") {
1879 let address = sub.put_blob(src.as_bytes())?;
1880 promoted.insert("code_address".into(), Value::from(address));
1881 }
1882 let mut spec = crate::substrate::GrainSpec::new(
1883 crate::model::grain_type::FACT,
1884 LOOP_NS,
1885 )
1886 .with_field("subject", rec.target_ref.clone())
1887 .with_field("relation", relation)
1888 .with_field(
1889 "object",
1890 serde_json::to_string(&promoted)
1891 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1892 )
1893 .with_field("rec_hash", rec_hash.to_string());
1894 if let Some(g) = gating {
1895 spec = spec
1896 .with_field("gating_evalset", g.evalset_hash.clone())
1897 .with_field("gating_run_id", g.run_id.clone());
1898 }
1899 created.push(sub.put_grain(&spec)?);
1900 }
1901 Proposal::Data { data } => {
1902 let revert_of = data
1907 .get("revert_of")
1908 .and_then(Value::as_str)
1909 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1910 self.rollback(
1911 sub,
1912 revert_of,
1913 actor,
1914 observer,
1915 scopes,
1916 &because,
1917 now_ms,
1918 )?;
1919 p = LoopPersisted::from_value(sub.load_state()?)?;
1922 if let Ok(reverted) = load_rec(sub, revert_of) {
1932 strike_cooldown(&mut p, reverted.dedup_key, now_ms);
1933 }
1934 }
1935 }
1936
1937 let applied = AppliedRecord {
1938 applied_at_ms: now_ms,
1939 target_ref: rec.target_ref.clone(),
1940 rollbackable: rec.rollbackable,
1941 created_hashes: created,
1942 inverse_cal,
1943 metric: rec.metric.clone(),
1944 };
1945 let prev = p.audit_heads.get(rec_hash).cloned();
1946 let audit = AuditRecord {
1947 rec_hash: rec_hash.into(),
1948 from: Some(status),
1949 to: RecStatus::Applied,
1950 actor: actor.into(),
1951 observer_type: observer,
1952 because,
1953 previous_audit_hash: prev,
1954 gating: gating.cloned(),
1955 at_ms: now_ms,
1956 };
1957 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1958 p.audit_heads.insert(rec_hash.into(), audit_hash);
1959 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1960 p.applied.insert(rec_hash.into(), applied.clone());
1961 sub.store_state(&p.to_value()?)?;
1962 Ok(applied)
1963 }
1964
1965 #[allow(clippy::too_many_arguments)]
1968 pub fn rollback<S: OmsSubstrate>(
1969 &self,
1970 sub: &mut S,
1971 rec_hash: &str,
1972 actor: &str,
1973 observer: ObserverType,
1974 scopes: &ScopeSet,
1975 because: &str,
1976 now_ms: i64,
1977 ) -> Result<()> {
1978 if !scopes.has(Scope::Apply) {
1979 return Err(Error::ScopeDenied("apply".into()));
1980 }
1981 let because = validate_because(because)?;
1982 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1983 let status = *p
1984 .status_index
1985 .get(rec_hash)
1986 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1987 if !status.can_transition_to(RecStatus::RolledBack, false) {
1988 return Err(Error::LifecycleViolation(format!(
1989 "{} -> rolled_back",
1990 status.as_str()
1991 )));
1992 }
1993 let applied = p
1994 .applied
1995 .get(rec_hash)
1996 .cloned()
1997 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1998 if !applied.rollbackable {
1999 return Err(Error::LifecycleViolation(
2000 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
2001 ));
2002 }
2003 for h in &applied.created_hashes {
2004 sub.retract(h, &format!("rollback of {rec_hash}"))?;
2005 }
2006 if let Some(inverse) = &applied.inverse_cal {
2013 sub.execute_cal(inverse)?;
2014 }
2015 let prev = p.audit_heads.get(rec_hash).cloned();
2016 let audit = AuditRecord {
2017 rec_hash: rec_hash.into(),
2018 from: Some(status),
2019 to: RecStatus::RolledBack,
2020 actor: actor.into(),
2021 observer_type: observer,
2022 because,
2023 previous_audit_hash: prev,
2024 gating: None,
2025 at_ms: now_ms,
2026 };
2027 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
2028 p.audit_heads.insert(rec_hash.into(), audit_hash);
2029 p.status_index
2030 .insert(rec_hash.into(), RecStatus::RolledBack);
2031 sub.store_state(&p.to_value()?)?;
2032 Ok(())
2033 }
2034
2035 pub fn recommendations<S: OmsSubstrate>(
2040 &self,
2041 sub: &S,
2042 status_filter: Option<RecStatus>,
2043 ) -> Result<Vec<Recommendation>> {
2044 let p = LoopPersisted::from_value(sub.load_state()?)?;
2045 let grains = sub.grains_of_type(
2046 crate::model::grain_type::RECOMMENDATION,
2047 Some(LOOP_NS),
2048 ReadOpts {
2049 live_only: false,
2050 since_ms: None,
2051 },
2052 )?;
2053 let mut out = Vec::new();
2054 for g in grains {
2055 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
2056 rec.status = p
2057 .status_index
2058 .get(&g.hash)
2059 .copied()
2060 .unwrap_or(RecStatus::Pending);
2061 if let Some(f) = status_filter {
2062 if rec.status != f {
2063 continue;
2064 }
2065 }
2066 out.push(rec);
2067 }
2068 out.sort_by(|a, b| {
2077 b.severity
2078 .cmp(&a.severity)
2079 .then(a.created_at_ms.cmp(&b.created_at_ms))
2080 .then(a.dedup_key.cmp(&b.dedup_key))
2081 .then(a.hash.cmp(&b.hash))
2082 });
2083 Ok(out)
2084 }
2085
2086 pub fn analyzer_settings<S: OmsSubstrate>(
2089 &self,
2090 sub: &S,
2091 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
2092 let p = LoopPersisted::from_value(sub.load_state()?)?;
2093 Ok(self
2094 .analyzers
2095 .iter()
2096 .map(|a| {
2097 let m = a.manifest();
2098 let cfg = p.config.get(&m.id);
2099 crate::config::AnalyzerSetting {
2100 id: m.id.clone(),
2101 title: m.title.clone(),
2102 description: m.description.clone(),
2103 tier: format!("{:?}", m.tier),
2104 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
2105 default_on: m.default_on,
2106 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
2107 severity_floor: cfg
2108 .and_then(|c| c.severity_floor)
2109 .map(|s| s.as_str().to_string()),
2110 }
2111 })
2112 .collect())
2113 }
2114
2115 pub fn set_analyzer_config<S: OmsSubstrate>(
2121 &self,
2122 sub: &mut S,
2123 analyzer_id: &str,
2124 update: crate::config::AnalyzerConfigUpdate,
2125 scopes: &ScopeSet,
2126 ) -> Result<crate::config::AnalyzerConfig> {
2127 if !scopes.has(Scope::Admin) {
2128 return Err(Error::ScopeDenied("admin".into()));
2129 }
2130 let manifest = self
2131 .analyzers
2132 .iter()
2133 .map(|a| a.manifest())
2134 .find(|m| m.id == analyzer_id)
2135 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
2136 if let Some(params) = &update.params {
2138 manifest.resolve_params(params)?;
2139 }
2140 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
2141 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
2142 if let Some(enabled) = update.enabled {
2143 cfg.enabled = Some(enabled);
2144 }
2145 if update.clear_floor {
2146 cfg.severity_floor = None;
2147 } else if let Some(floor) = update.severity_floor {
2148 cfg.severity_floor = Some(floor);
2149 }
2150 if let Some(params) = update.params {
2151 cfg.params = params;
2152 }
2153 if let Some(ns) = update.namespaces {
2154 cfg.namespaces = ns;
2155 }
2156 let stored = cfg.clone();
2157 sub.store_state(&p.to_value()?)?;
2158 Ok(stored)
2159 }
2160
2161 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
2164 let p = LoopPersisted::from_value(sub.load_state()?)?;
2165 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
2166 out.sort_by(|a, b| {
2170 a.measured_at_ms
2171 .cmp(&b.measured_at_ms)
2172 .then(a.horizon_ms.cmp(&b.horizon_ms))
2173 .then(a.metric.cmp(&b.metric))
2174 .then(a.rec_hash.cmp(&b.rec_hash))
2175 });
2176 Ok(out)
2177 }
2178
2179 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
2183 let p = LoopPersisted::from_value(sub.load_state()?)?;
2184 let new = count_new(sub, p.state.watermark_ms)?;
2185 let (grains_since_run, error_events_since_run) = (new.grains, new.error_events);
2186 let recs = self.recommendations(sub, None)?;
2187 let mut pending = 0;
2188 let mut applied = 0;
2189 for r in &recs {
2190 match r.status {
2191 RecStatus::Pending => pending += 1,
2192 RecStatus::Applied => applied += 1,
2193 _ => {}
2194 }
2195 }
2196 let stale = match p.state.last_run_ms {
2198 None => true,
2199 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
2200 };
2201 Ok(Health {
2202 last_run_ms: p.state.last_run_ms,
2203 grains_since_run,
2204 error_events_since_run,
2205 pending,
2206 applied,
2207 total: recs.len() as u64,
2208 stale,
2209 })
2210 }
2211
2212 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
2217 let recs = self.recommendations(sub, None)?;
2218 let mut m = LlmMetrics::default();
2219 for r in &recs {
2220 if !matches!(r.origin, Origin::Llm { .. }) {
2221 continue;
2222 }
2223 m.proposed += 1;
2224 match r.status {
2225 RecStatus::Pending => m.pending += 1,
2226 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
2227 RecStatus::Rejected => m.rejected += 1,
2228 RecStatus::Expired | RecStatus::Withdrawn => {}
2231 }
2232 }
2233 let decided = m.approved + m.rejected;
2234 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
2235 Ok(m)
2236 }
2237}
2238
2239#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
2241pub struct Health {
2242 #[serde(skip_serializing_if = "Option::is_none")]
2243 pub last_run_ms: Option<i64>,
2244 pub grains_since_run: u64,
2245 pub error_events_since_run: u64,
2246 pub pending: u64,
2247 pub applied: u64,
2248 pub total: u64,
2249 pub stale: bool,
2252}
2253
2254#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
2256pub struct LlmMetrics {
2257 pub proposed: u64,
2260 pub pending: u64,
2261 pub approved: u64,
2263 pub rejected: u64,
2264 #[serde(skip_serializing_if = "Option::is_none")]
2266 pub approval_rate: Option<f64>,
2267}
2268
2269fn measure_outcomes<S: OmsSubstrate>(
2276 sub: &S,
2277 p: &mut LoopPersisted,
2278 policy: &crate::policy::Policy,
2279 now_ms: i64,
2280) -> Result<Vec<OutcomeInput>> {
2281 let mut due: Vec<(String, crate::config::AppliedRecord, Checkpoint)> = Vec::new();
2285 for (h, a) in &p.applied {
2286 if p.status_index.get(h) != Some(&RecStatus::Applied) {
2287 continue;
2288 }
2289 let Some(metric) = &a.metric else { continue };
2290 let done = p.measured.get(h).cloned().unwrap_or_default();
2291 for cp in metric.schedule() {
2292 if !done.contains(&cp) && checkpoint_due(sub, metric, a.applied_at_ms, cp, now_ms)? {
2293 due.push((h.clone(), a.clone(), cp));
2294 }
2295 }
2296 }
2297
2298 let mut out = Vec::new();
2299 for (rec_hash, applied, checkpoint) in due {
2300 let metric = applied.metric.as_ref().unwrap();
2301 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
2302 continue; };
2304 let bound = cost_bound_for(policy, metric);
2305 let base = baseline_at_apply(
2306 sub,
2307 metric,
2308 applied.applied_at_ms,
2309 baseline_kind_for(policy, metric),
2310 bound.map(|b| b.field.as_str()),
2311 )?;
2312 let baseline = base.value;
2313 let tolerance = tolerance_for(policy, metric, &base);
2316 let regressed = crate::recommendation::is_regression(
2317 baseline,
2318 current,
2319 metric.higher_is_better,
2320 tolerance,
2321 );
2322 let (current_run_id, cost) = current_run_and_cost(sub, metric, applied.applied_at_ms, bound, &base)?;
2327 let costlier = cost.as_ref().is_some_and(|c| c.breached());
2328 let verdict = if regressed {
2329 "regressed"
2330 } else if costlier {
2331 "held_costlier"
2332 } else {
2333 "held"
2334 };
2335 p.outcomes.entry(rec_hash.clone()).or_default().push(
2336 crate::recommendation::OutcomeResult {
2337 rec_hash: rec_hash.clone(),
2338 metric: metric.metric.clone(),
2339 baseline,
2340 current,
2341 verdict: verdict.into(),
2342 baseline_kind: base.kind.into(),
2343 baseline_run_id: base.run_id.clone(),
2344 best_before: base.best_before,
2345 tolerance,
2346 current_run_id: current_run_id.clone(),
2347 cost: cost.clone(),
2348 horizon_ms: checkpoint.as_ms().unwrap_or(0),
2351 checkpoint: checkpoint.as_ms().is_none().then_some(checkpoint),
2352 measured_at_ms: now_ms,
2353 },
2354 );
2355 p.measured.entry(rec_hash.clone()).or_default().push(checkpoint);
2356 if regressed || costlier {
2357 out.push(OutcomeInput {
2358 rec_hash,
2359 target_ref: applied.target_ref.clone(),
2360 metric: metric.metric.clone(),
2361 baseline,
2362 current,
2363 unit: metric.unit.clone(),
2364 higher_is_better: metric.higher_is_better,
2365 baseline_kind: base.kind.into(),
2366 baseline_run_id: base.run_id,
2367 best_before: base.best_before,
2368 tolerance,
2369 current_run_id,
2370 cost,
2371 });
2372 }
2373 }
2374 Ok(out)
2375}
2376
2377fn cost_bound_for<'p>(
2379 policy: &'p crate::policy::Policy,
2380 metric: &crate::recommendation::MetricSnapshot,
2381) -> Option<&'p crate::policy::CostBound> {
2382 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2383 (Some(e), Some((hash, _))) if e.hash == hash => e.cost.as_ref(),
2384 _ => None,
2385 }
2386}
2387
2388fn current_run_and_cost<S: SubstrateRead>(
2394 sub: &S,
2395 metric: &crate::recommendation::MetricSnapshot,
2396 applied_at_ms: i64,
2397 bound: Option<&crate::policy::CostBound>,
2398 base: &BaselineRead,
2399) -> Result<(Option<String>, Option<crate::recommendation::CostRead>)> {
2400 let Some((evalset, _)) = crate::eval::parse_evalset_metric(&metric.metric) else {
2401 return Ok((None, None));
2402 };
2403 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(applied_at_ms))? else {
2404 return Ok((None, None));
2405 };
2406 let cost = bound.map(|b| {
2407 let current = crate::eval::run_value(&run, &b.field);
2408 let status = match (base.cost, current) {
2409 (Some(bl), Some(cur)) if cur > bl * b.max_increase_ratio + 1e-9 => "breached",
2410 (Some(_), Some(_)) => "within",
2411 _ => "not_measurable",
2412 };
2413 crate::recommendation::CostRead {
2414 field: b.field.clone(),
2415 max_increase_ratio: b.max_increase_ratio,
2416 baseline: base.cost,
2417 current,
2418 status: status.into(),
2419 }
2420 });
2421 Ok((Some(run.run_id), cost))
2422}
2423
2424fn tolerance_for(
2429 policy: &crate::policy::Policy,
2430 metric: &crate::recommendation::MetricSnapshot,
2431 base: &BaselineRead,
2432) -> f64 {
2433 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2434 (Some(e), Some((hash, field))) if e.hash == hash => e
2435 .min_effect
2436 .map(|m| m.resolve(field, base.total.unwrap_or(metric.n)))
2437 .unwrap_or(0.0),
2438 _ => 0.0,
2439 }
2440}
2441
2442fn checkpoint_due<S: SubstrateRead>(
2451 sub: &S,
2452 metric: &crate::recommendation::MetricSnapshot,
2453 applied_at_ms: i64,
2454 cp: Checkpoint,
2455 now_ms: i64,
2456) -> Result<bool> {
2457 Ok(match cp {
2458 Checkpoint::AfterMs(ms) => now_ms - applied_at_ms >= ms,
2459 Checkpoint::AfterRuns(n) => match crate::eval::parse_evalset_metric(&metric.metric) {
2460 Some((evalset, _)) => {
2461 crate::eval::eval_runs(sub, evalset, Some(applied_at_ms + 1))?.len() >= n as usize
2462 }
2463 None => false,
2464 },
2465 Checkpoint::AfterGrains(n) => count_new(sub, Some(applied_at_ms))?.grains >= n as u64,
2466 })
2467}
2468
2469pub(crate) struct BaselineRead {
2471 pub value: f64,
2472 pub kind: &'static str,
2475 pub run_id: Option<String>,
2476 pub best_before: Option<f64>,
2477 pub total: Option<u64>,
2479 pub cost: Option<f64>,
2482}
2483
2484fn baseline_kind_for(
2487 policy: &crate::policy::Policy,
2488 metric: &crate::recommendation::MetricSnapshot,
2489) -> crate::policy::BaselineKind {
2490 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2491 (Some(e), Some((hash, _))) if e.hash == hash => e.baseline,
2492 _ => crate::policy::BaselineKind::default(),
2493 }
2494}
2495
2496pub(crate) fn baseline_at_apply<S: SubstrateRead>(
2514 sub: &S,
2515 metric: &crate::recommendation::MetricSnapshot,
2516 applied_at_ms: i64,
2517 kind: crate::policy::BaselineKind,
2518 cost_field: Option<&str>,
2519) -> Result<BaselineRead> {
2520 use crate::policy::BaselineKind;
2521 if let Some((evalset, field)) = crate::eval::parse_evalset_metric(&metric.metric) {
2522 let before: Vec<(crate::eval::EvalRun, f64)> = crate::eval::eval_runs(sub, evalset, None)?
2525 .into_iter()
2526 .filter(|r| r.recorded_ms < applied_at_ms)
2527 .filter_map(|r| crate::eval::run_value(&r, field).map(|v| (r, v)))
2528 .collect();
2529 if let Some(newest) = before.last() {
2530 let best = before
2533 .iter()
2534 .fold(None::<&(crate::eval::EvalRun, f64)>, |acc, r| match acc {
2535 None => Some(r),
2536 Some(b) => {
2537 let better = if metric.higher_is_better { r.1 > b.1 } else { r.1 < b.1 };
2538 Some(if better { r } else { b })
2539 }
2540 })
2541 .expect("non-empty");
2542 let pick = match kind {
2543 BaselineKind::NewestBeforeApply => newest,
2544 BaselineKind::HighWater => best,
2545 };
2546 return Ok(BaselineRead {
2547 value: pick.1,
2548 kind: kind.as_str(),
2549 run_id: Some(pick.0.run_id.clone()),
2550 best_before: Some(best.1),
2551 total: Some(pick.0.total()),
2552 cost: cost_field.and_then(|f| crate::eval::run_value(&pick.0, f)),
2553 });
2554 }
2555 }
2556 Ok(BaselineRead { value: metric.baseline, kind: "snapshot", run_id: None, best_before: None, total: None, cost: None })
2557}
2558
2559pub(crate) fn measure_metric<S: SubstrateRead>(
2561 sub: &S,
2562 metric: &crate::recommendation::MetricSnapshot,
2563 since_ms: i64,
2564) -> Result<Option<f64>> {
2565 match metric.metric.as_str() {
2566 "tool_error_recurrence" => {
2571 let Some(tool) = &metric.subject else { return Ok(None) };
2572 let tools = sub.grains_of_type(
2573 crate::model::grain_type::TOOL,
2574 None,
2575 ReadOpts { live_only: true, since_ms: Some(since_ms) },
2576 )?;
2577 let n = tools
2578 .iter()
2579 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
2580 .filter(|t| {
2581 metric.relation.as_deref().is_none_or(|sig| {
2584 crate::analyzers::tool_failure::normalize_signature(
2585 t.tool_content().unwrap_or(""),
2586 ) == sig
2587 })
2588 })
2589 .count();
2590 Ok(Some(n as f64))
2591 }
2592 "contradiction_recurrence" => {
2596 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
2597 return Ok(None);
2598 };
2599 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
2600 let distinct: BTreeSet<String> = facts
2601 .iter()
2602 .filter(|f| {
2603 f.fact_relation()
2604 .is_some_and(|r| normalize_ident(r) == *relation)
2605 })
2606 .filter_map(|f| f.fact_object().map(normalize_ident))
2607 .collect();
2608 Ok(Some(distinct.len().saturating_sub(1) as f64))
2609 }
2610 m if m.starts_with("evalset:") => {
2623 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
2624 return Ok(None);
2625 };
2626 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
2627 return Ok(None);
2628 };
2629 Ok(crate::eval::run_value(&run, field))
2630 }
2631 _ => Ok(None),
2632 }
2633}
2634
2635fn scoped_live_facts<S: SubstrateRead>(
2638 sub: &S,
2639 namespace: Option<&str>,
2640 subject: &str,
2641) -> Result<Vec<GrainRecord>> {
2642 let facts = sub.grains_of_type(
2643 crate::model::grain_type::FACT,
2644 None,
2645 ReadOpts { live_only: true, since_ms: None },
2646 )?;
2647 Ok(facts
2648 .into_iter()
2649 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
2650 .filter(|f| {
2651 f.fact_subject()
2652 .is_some_and(|s| normalize_ident(s) == subject)
2653 })
2654 .collect())
2655}
2656
2657fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
2669 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
2670 return false;
2671 };
2672 if fields.is_empty() {
2673 return false;
2674 }
2675 let Ok(Some(grain)) = sub.grain(&target) else {
2676 return false;
2677 };
2678 if grain.valid_to_ms.is_some() {
2688 return false;
2689 }
2690 fields.iter().all(|(k, v)| {
2692 if k == "namespace" {
2693 return v
2694 .as_str()
2695 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2696 }
2697 match (v, grain.fields.get(k)) {
2698 (Value::String(a), Some(Value::String(b))) => {
2699 normalize_ident(a) == normalize_ident(b)
2700 }
2701 (a, Some(b)) => a == b,
2702 (_, None) => false,
2703 }
2704 })
2705}
2706
2707const EVIDENCE_CAP: usize = 64;
2717const CITED_SEED_CAP: usize = 24;
2718const TOOL_SEED_CAP: usize = 16;
2719const NOTE_SEED_CAP: usize = 8;
2724const HARNESS_EVIDENCE_KINDS: &[&str] = &["fold_summary"];
2743const HARNESS_SEED_CAP: usize = 6;
2744const LENS_RESERVE: usize = 24;
2745
2746const MIN_LLM_CONFIDENCE: f64 = 0.75;
2749
2750macro_rules! discover_instructions {
2760 ($scoring:literal) => {
2761 concat!(
2762 "You review an agent's memory for quality. \
2763Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2764findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2765stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2766recurring cost or hand-off the agent's own setup could remove). \
2767The deterministic findings already cover what the ERROR TEXT says; restating \
2768one of them earns nothing. The evidence may also contain OUTCOME records — a \
2769run's observable shape together with whether it was accepted or rejected. A \
2770problem that raised no error at all is exactly the kind the deterministic \
2771checks cannot see, so compare the rejected outcomes against the accepted \
2772ones: a feature they share and the accepted ones lack is a candidate rule. \
2773Require at least two rejected outcomes before proposing one — a single \
2774rejection is an anecdote, not a pattern. ",
2775 $scoring,
2776 " The 'approved' and 'rejected' lists, when \
2777present, show findings this reviewer recently accepted or rejected — prefer the \
2778kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2779or more evidence items from the bundle by their 'id' (or 'hash'), name a \
2780'target', and include your confidence 0.0-1.0. Return JSON: \
2781{\"recommendations\":[{\"summary\":\"...\",\
2782\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<id>\"],\
2783\"confidence\":0.0,\"proposal\":{...}}]}. \
2784OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2785you are not asking to change anything. Include it ONLY when the evidence \
2786supports a specific change, choosing exactly one kind: \
2787(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2788\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
2789an action the agent itself takes on the next occasion. Either ADD an action \
2790it is failing to take ('Record the vendor name and the amount on every \
2791invoice, not just the date') or ORDER one that goes wrong ('Refund a \
2792subscription before cancelling it; refunds on cancelled subscriptions are \
2793refused'). Name the action, not a check on it: 'validate', 'verify' and \
2794'ensure ... is correct' describe a review step the agent has no way to \
2795perform, and such a rule changes nothing even once applied. \
2796(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2797entity target — a durable fact the agent keeps having to be told (an alias, a \
2798settled default, a preference). 'relation' is a short identifier (letters, \
2799digits, _ - . :), not a sentence. \
2800(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2801\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2802assembles the agent's context, when the evidence shows it retrieves the wrong \
2803things. Give the FULL new body; it replaces the old one. \
2804(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2805with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2806workflow. Only these paths are editable: 'edges.<i>.cond', \
2807'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2808holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2809(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2810— full replacement source for that tool. It is applied only after a recorded \
2811evaluation run passes, so propose one only when the evidence shows the current \
2812code is the defect. \
2813The subject of a fact, the name of a query, the plan hash and the tool name \
2814all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2815becomes a change a human reviewer may apply, so it must be fully supported by \
2816the cited evidence. Propose nothing you cannot ground in the evidence."
2817 )
2818 };
2819}
2820
2821const DISCOVER_INSTRUCTIONS: &str = discover_instructions!(
2823 "SCORING: propose a finding ONLY if you \
2824are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2825useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2826earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2827when there is nothing worth flagging."
2828);
2829
2830const DISCOVER_LEARNER_INSTRUCTIONS: &str = discover_instructions!(
2832 "SCORING: you are the learning stage of a deployed agent, and what you \
2833propose now is what it will do differently next time — a lesson you withhold \
2834is a mistake it repeats. A correct, actionable proposal earns 1; a wrong or \
2835trivial one is penalized 1; returning nothing while the evidence holds a \
2836recurring failure, two or more rejected outcomes, an instruction from a \
2837person, or a multi-step procedure the agent completed successfully that no \
2838saved skill or plan covers, is ALSO penalized 1. Abstain only when the \
2839evidence shows none of those. Prefer the one proposal that addresses the most \
2840frequent or most costly failure — or, when nothing failed, the procedure that \
2841worked — over several speculative ones, and report your confidence honestly — \
2842an independent verifier, not you, decides what survives."
2843);
2844
2845const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2851fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2852the facts it relies on are actually present in the cited evidence, NOT that its \
2853conclusion is stated verbatim. Decompose the finding into the factual claims it \
2854depends on. Mark supported=true when those facts are present in the evidence \
2855(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2856on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2857different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2858
2859const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2861each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2862never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2863SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2864'possible' findings with no concrete defect, and reject any claimed \
2865inconsistency or contradiction that is not backed by at least two actually \
2866conflicting facts in the cited evidence. (2) Context — does the finding \
2867correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2868finding when it names a genuine, specific problem grounded in its evidence and \
2869materially useful to a human reviewer; otherwise reject it, and default to \
2870keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2871redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2872grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2873KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2874{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2875
2876const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2878guidance note to help a human reviewer decide. Do not restate the finding. Return \
2879JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2880
2881fn push_evidence(
2887 evidence: &mut Vec<crate::llm::EvidenceItem>,
2888 bundle: &mut BTreeSet<String>,
2889 ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2890 g: &GrainRecord,
2891 attribution: crate::policy::EvidenceAttribution,
2892) {
2893 if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2894 ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2895 evidence.push(crate::llm::EvidenceItem {
2896 id: format!("e{}", evidence.len() + 1),
2897 hash: g.hash.clone(),
2898 grain_type: g.grain_type.clone(),
2899 text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
2900 });
2901 }
2902}
2903
2904pub(crate) fn resolve_citation(
2912 cite: &str,
2913 bundle: &BTreeSet<String>,
2914 id_to_hash: &std::collections::BTreeMap<&str, &str>,
2915) -> Option<String> {
2916 let cite = cite.trim();
2917 if bundle.contains(cite) {
2918 return Some(cite.to_string());
2919 }
2920 if let Some(h) = id_to_hash.get(cite) {
2921 return Some((*h).to_string());
2922 }
2923 const MIN_PREFIX: usize = 12;
2924 if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
2925 let lower = cite.to_ascii_lowercase();
2926 let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
2927 if let (Some(h), None) = (it.next(), it.next()) {
2928 return Some(h.clone());
2929 }
2930 }
2931 None
2932}
2933
2934fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
2937 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2938 return format!("{s} {r} {o}");
2939 }
2940 if let Some(t) = g.tool_name() {
2945 let status = if g.is_error() { "error" } else { "ok" };
2946 let out = g.tool_content().unwrap_or("");
2947 let input = match g.fields.get("input") {
2953 Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
2954 Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
2955 _ => String::new(),
2956 };
2957 return format!("tool {t}{input} {status}: {out}");
2958 }
2959 for key in ["content", "body", "text", "summary", "object"] {
2967 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2968 if v.is_empty() {
2969 continue;
2970 }
2971 if attribution == crate::policy::EvidenceAttribution::Anonymous {
2982 return v.to_string();
2983 }
2984 if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
2985 if !who.is_empty() {
2986 let kind = g
2987 .fields
2988 .get("observer_type")
2989 .and_then(|v| v.as_str())
2990 .unwrap_or("");
2991 let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
2992 let mut prefix = if kind == "human" {
2993 format!("{who} (a person) said")
2994 } else {
2995 format!("{who} observed")
2996 };
2997 if !about.is_empty() {
2998 prefix.push_str(&format!(" of {about}"));
2999 }
3000 return format!("{prefix}: {v}");
3001 }
3002 }
3003 return v.to_string();
3004 }
3005 }
3006 String::new()
3007}
3008
3009fn sanitize_lesson(s: &str) -> String {
3014 sanitize_line(s, crate::llm::MAX_LESSON_LEN)
3015}
3016
3017fn sanitize_line(s: &str, max: usize) -> String {
3023 let cleaned: String =
3024 s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
3025 crate::llm::cap(cleaned.trim(), max)
3026}
3027
3028fn sanitize_relation(s: &str) -> Option<String> {
3032 let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
3033 if r.is_empty()
3034 || !r
3035 .chars()
3036 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
3037 {
3038 return None;
3039 }
3040 Some(r)
3041}
3042
3043fn safe_definition_body(body: &str) -> bool {
3059 if body.contains('{') || body.contains('}') {
3060 return false;
3061 }
3062 !body
3063 .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
3064 .any(|tok| {
3065 ["FORGET", "PURGE", "DROP", "DEFINE"]
3066 .iter()
3067 .any(|kw| tok.eq_ignore_ascii_case(kw))
3068 })
3069}
3070
3071fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
3076 let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3077 match resolved {
3078 Some(r) => format!("{summary} {}", r.rendered),
3079 None => summary,
3080 }
3081}
3082
3083struct DecidedDraft {
3090 grounded: bool,
3091 sound: f64,
3093 judged_by: crate::decide::JudgedBy,
3094}
3095
3096struct ValidatedDraft {
3097 draft: crate::llm::LlmDraft,
3098 target_ref: String,
3099 cited: Vec<String>,
3100 resolved: Option<ResolvedProposal>,
3101}
3102
3103struct ResolvedProposal {
3107 action: ActionKind,
3108 proposal: Proposal,
3109 rendered: String,
3112 summary_key: &'static str,
3113 summary_args: serde_json::Map<String, Value>,
3114 rollbackable: bool,
3115 evalset_hash: Option<String>,
3116 importance: f64,
3117 fact_fields: Option<serde_json::Map<String, Value>>,
3121 extra_statements: Vec<String>,
3124 replay: Option<Value>,
3127}
3128
3129fn plan_edit_allowed(path: &str) -> bool {
3139 let seg: Vec<&str> = path.split('.').collect();
3140 match seg.as_slice() {
3141 ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
3142 ["retries", node] => !node.is_empty(),
3143 _ => false,
3144 }
3145}
3146
3147fn plan_get(body: &Value, path: &str) -> Value {
3150 let mut cur = body;
3151 for seg in path.split('.') {
3152 cur = match cur {
3153 Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
3154 Some(v) => v,
3155 None => return Value::Null,
3156 },
3157 Value::Object(o) => match o.get(seg) {
3158 Some(v) => v,
3159 None => return Value::Null,
3160 },
3161 _ => return Value::Null,
3162 };
3163 }
3164 cur.clone()
3165}
3166
3167fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
3173 let segs: Vec<&str> = path.split('.').collect();
3174 let Some((last, parents)) = segs.split_last() else {
3175 return false;
3176 };
3177 let mut cur = body;
3178 for (depth, seg) in parents.iter().enumerate() {
3179 cur = match cur {
3180 Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
3181 Some(v) => v,
3182 None => return false,
3183 },
3184 Value::Object(o) => {
3185 if depth == 0 && *seg == "retries" && !o.contains_key("retries") {
3186 o.insert("retries".into(), Value::Object(serde_json::Map::new()));
3187 }
3188 match o.get_mut(*seg) {
3189 Some(v) => v,
3190 None => return false,
3191 }
3192 }
3193 _ => return false,
3194 };
3195 }
3196 match cur {
3197 Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
3198 Some(slot) => {
3199 *slot = to;
3200 true
3201 }
3202 None => false,
3203 },
3204 Value::Object(o) => {
3205 o.insert((*last).to_string(), to);
3206 true
3207 }
3208 _ => false,
3209 }
3210}
3211
3212fn plan_value_ok(path: &str, to: &Value) -> bool {
3217 let seg: Vec<&str> = path.split('.').collect();
3218 match seg.as_slice() {
3219 ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
3220 !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
3221 }),
3222 ["edges", _, "max_cycles"] | ["retries", _] => {
3223 to.as_u64().is_some_and(|n| n <= 1_000)
3224 }
3225 _ => false,
3226 }
3227}
3228
3229fn resolve_proposal<S: OmsSubstrate>(
3235 sub: &S,
3236 d: &crate::llm::LlmDraft,
3237 target: &TargetRef,
3238 cited: &[String],
3239 ns_by_hash: &std::collections::BTreeMap<String, String>,
3240 caps: Capabilities,
3241 policy: &crate::policy::Policy,
3242) -> Option<ResolvedProposal> {
3243 use crate::llm::DraftProposal as P;
3244 let (skills, plans) = (&policy.skills, &policy.plans);
3245 let mut args = serde_json::Map::new();
3246 match d.parsed_proposal()? {
3247 P::Plan { description, when_to_use, nodes, edges } => {
3249 if !plans.enabled || !caps.plans {
3250 return None;
3251 }
3252 let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
3253 derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
3254 args.insert("name".into(), Value::from(name.clone()));
3255 args.insert("nodes".into(), Value::from(n_nodes as u64));
3256 let mut stmts = vec![match &existing_skill {
3257 Some(h) => cal::supersede(h, "skill", &skill),
3258 None => cal::add("skill", &skill),
3259 }];
3260 let (summary_key, kind) = match &workflow {
3267 Some(wf) => {
3268 stmts.push(match &existing_plan {
3269 Some(h) => cal::supersede(h, "workflow", wf),
3270 None => cal::add("workflow", wf),
3271 });
3272 args.insert("edges".into(), Value::from(n_edges as u64));
3273 ("llm.plan", "plan")
3274 }
3275 None => {
3276 args.insert("steps".into(), Value::from(n_nodes as u64));
3277 ("llm.skill", "skill")
3278 }
3279 };
3280 let patched = existing_skill.is_some() || (workflow.is_some() && existing_plan.is_some());
3281 let (action, verb) = if patched {
3282 (ActionKind::Revise, "revise")
3283 } else {
3284 (ActionKind::Record, "record")
3285 };
3286 Some(ResolvedProposal {
3287 action,
3288 proposal: Proposal::Cal { cal: cal::batch(&stmts) },
3289 rendered: format!(
3290 "Proposed {kind} to {verb}: \"{name}\" — {n_nodes} steps; when: {}",
3291 skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
3292 ),
3293 summary_key,
3294 summary_args: args,
3295 rollbackable: true,
3296 evalset_hash: None,
3297 importance: 0.65,
3298 fact_fields: None,
3299 extra_statements: Vec::new(),
3300 replay: None,
3301 })
3302 }
3303 P::Skill { description, when_to_use, steps } => {
3305 if !skills.enabled {
3306 return None;
3307 }
3308 let SkillFields { fields, name, n_steps, existing } =
3309 derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
3310 args.insert("name".into(), Value::from(name.clone()));
3311 args.insert("steps".into(), Value::from(n_steps as u64));
3312 let (action, cal, verb) = match existing {
3315 Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
3316 None => (ActionKind::Record, cal::add("skill", &fields), "record"),
3317 };
3318 Some(ResolvedProposal {
3319 action,
3320 proposal: Proposal::Cal { cal },
3321 rendered: format!(
3322 "Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
3323 fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
3324 ),
3325 summary_key: "llm.skill",
3326 summary_args: args,
3327 rollbackable: true,
3328 evalset_hash: None,
3329 importance: 0.6,
3330 fact_fields: None,
3331 extra_statements: Vec::new(),
3332 replay: None,
3333 })
3334 }
3335 P::Consolidation { lesson, supersedes } => {
3337 let lesson = sanitize_lesson(&lesson);
3338 if lesson.is_empty() {
3339 return None;
3340 }
3341 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
3342 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("").to_string();
3343 let mut hashes: Vec<String> = supersedes.into_iter().collect();
3348 hashes.sort();
3349 hashes.dedup();
3350 let mut members = Vec::new();
3351 for h in &hashes {
3352 let g = sub.grain(h).ok().flatten()?;
3353 if !g.is_live()
3354 || g.fact_relation() != Some("lesson")
3355 || g.fact_subject().is_none_or(|s| normalize_ident(s) != normalize_ident(&subject))
3356 {
3357 return None;
3358 }
3359 members.push(g);
3360 }
3361 if members.len() < 2 || members.iter().any(|m| m.namespace != members[0].namespace) {
3362 return None;
3363 }
3364 let ns = members[0].namespace.clone();
3365 let mut fields = fields;
3366 if !ns.is_empty() {
3367 fields.insert("namespace".into(), Value::from(ns.clone()));
3368 }
3369 fields.insert("consolidates".into(), Value::from(hashes.clone()));
3370 let extra_statements: Vec<String> = members
3375 .iter()
3376 .map(|m| {
3377 let mut marker = serde_json::Map::new();
3378 marker.insert("subject".into(), Value::from(subject.clone()));
3379 marker.insert("relation".into(), Value::from("mg:lesson_consolidated"));
3380 marker.insert("object".into(), Value::from(lesson.clone()));
3381 if !ns.is_empty() {
3382 marker.insert("namespace".into(), Value::from(ns.clone()));
3383 }
3384 cal::supersede(&m.hash, "fact", &marker)
3385 })
3386 .collect();
3387 args.insert("lesson".into(), Value::from(lesson.clone()));
3388 args.insert("count".into(), Value::from(members.len() as u64));
3389 Some(ResolvedProposal {
3390 action: ActionKind::Consolidate,
3391 proposal: Proposal::Cal { cal: cal::batch(&extra_statements) },
3392 rendered: format!(
3393 "Proposed consolidation of {} lessons on \"{subject}\" into one: \"{lesson}\"",
3394 members.len()
3395 ),
3396 summary_key: "llm.consolidation",
3397 summary_args: args,
3398 rollbackable: true,
3399 evalset_hash: None,
3400 importance: 0.6,
3401 fact_fields: Some(fields),
3402 extra_statements,
3403 replay: None,
3404 })
3405 }
3406 P::Lesson { lesson } => {
3408 let lesson = sanitize_lesson(&lesson);
3409 if lesson.is_empty() {
3410 return None;
3411 }
3412 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
3413 args.insert("lesson".into(), Value::from(lesson.clone()));
3414 Some(ResolvedProposal {
3415 action: ActionKind::ClusterFailure,
3419 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
3420 rendered: format!("Proposed lesson to record: \"{lesson}\""),
3421 summary_key: "llm.lesson",
3422 summary_args: args,
3423 rollbackable: true,
3424 evalset_hash: None,
3425 importance: 0.5,
3426 fact_fields: Some(fields),
3427 extra_statements: Vec::new(),
3428 replay: None,
3429 })
3430 }
3431 P::Fact { relation, object } => {
3433 let relation = sanitize_relation(&relation)?;
3434 let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
3435 if object.is_empty() {
3436 return None;
3437 }
3438 let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
3439 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
3440 args.insert("relation".into(), Value::from(relation.clone()));
3441 args.insert("object".into(), Value::from(object.clone()));
3442 Some(ResolvedProposal {
3443 action: ActionKind::Record,
3444 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
3445 rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
3446 summary_key: "llm.fact",
3447 summary_args: args,
3448 rollbackable: true,
3449 evalset_hash: None,
3450 importance: 0.5,
3451 fact_fields: Some(fields),
3452 extra_statements: Vec::new(),
3453 replay: None,
3454 })
3455 }
3456 P::QueryRevision { body } => {
3458 let name = target.opaque();
3459 if name.is_empty()
3463 || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
3464 {
3465 return None;
3466 }
3467 let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
3468 if body.is_empty() || !safe_definition_body(&body) {
3469 return None;
3470 }
3471 let stmt = match target.scheme() {
3472 "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
3473 "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
3474 _ => return None,
3475 };
3476 sub.validate_cal(&stmt).ok()?;
3481 sub.definition_inverse(&stmt).ok().flatten()?;
3482 args.insert("name".into(), Value::from(name));
3483 args.insert("body".into(), Value::from(body.clone()));
3484 Some(ResolvedProposal {
3485 action: ActionKind::Revise,
3486 proposal: Proposal::Cal { cal: stmt },
3487 rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
3488 summary_key: "llm.query_revision",
3489 summary_args: args,
3490 rollbackable: true,
3491 evalset_hash: None,
3492 importance: 0.6,
3493 fact_fields: None,
3494 extra_statements: Vec::new(),
3495 replay: None,
3496 })
3497 }
3498 P::PlanRevision { edits } => {
3500 if !caps.plans
3501 || target.scheme() != "grain"
3502 || edits.is_empty()
3503 || edits.len() > crate::llm::MAX_PLAN_EDITS
3504 {
3505 return None;
3506 }
3507 let hash = target.opaque();
3508 let g = sub.grain(hash).ok().flatten()?;
3509 if g.grain_type != "workflow" || !g.is_live() {
3510 return None;
3511 }
3512 let mut body = Value::Object(g.fields.clone());
3513 let mut deltas = Vec::new();
3514 let nodes: std::collections::BTreeSet<String> = body
3515 .get("nodes")
3516 .and_then(Value::as_array)
3517 .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
3518 .unwrap_or_default();
3519 for e in &edits {
3520 if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
3521 return None;
3522 }
3523 if let Some(node) = e.path.strip_prefix("retries.") {
3528 if !nodes.contains(node) {
3529 return None;
3530 }
3531 }
3532 if plan_get(&body, &e.path) != e.from {
3535 return None;
3536 }
3537 if e.from == e.to {
3541 return None;
3542 }
3543 if !plan_set(&mut body, &e.path, e.to.clone()) {
3544 return None;
3545 }
3546 deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
3547 }
3548 sub.validate_plan(&body).ok()?;
3552 let replay = sub.plan_replay(hash, &body).ok().flatten();
3561 let refused = match (&policy.plan_replay, &replay) {
3562 (Some(gate), Some(report)) => gate.refusal(report),
3563 _ => None,
3564 };
3565 let Value::Object(fields) = body else {
3566 return None;
3567 };
3568 args.insert("plan".into(), Value::from(hash));
3569 args.insert("edits".into(), Value::from(deltas.join("; ")));
3570 if let Some(reason) = refused {
3571 let mut data = serde_json::Map::new();
3572 data.insert("plan".into(), Value::from(hash));
3573 data.insert("edits".into(), Value::from(deltas.clone()));
3574 data.insert("refused".into(), Value::from(reason.clone()));
3575 args.insert("reason".into(), Value::from(reason.clone()));
3576 return Some(ResolvedProposal {
3577 action: ActionKind::Flag,
3578 proposal: Proposal::Data { data },
3579 rendered: format!(
3580 "Plan revision ({}) refused by the rehearsal: {reason}",
3581 deltas.join("; ")
3582 ),
3583 summary_key: "llm.plan_revision_refused",
3584 summary_args: args,
3585 rollbackable: false,
3586 evalset_hash: None,
3587 importance: 0.4,
3588 fact_fields: None,
3589 extra_statements: Vec::new(),
3590 replay,
3591 });
3592 }
3593 let stmt = cal::supersede(hash, "workflow", &fields);
3594 sub.validate_cal(&stmt).ok()?;
3599 Some(ResolvedProposal {
3600 action: ActionKind::Revise,
3601 proposal: Proposal::Cal { cal: stmt },
3602 rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
3603 summary_key: "llm.plan_revision",
3604 summary_args: args,
3605 rollbackable: true,
3606 evalset_hash: None,
3607 importance: 0.7,
3608 fact_fields: None,
3609 extra_statements: Vec::new(),
3610 replay,
3611 })
3612 }
3613 P::CodeRevision { source } => {
3615 if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
3616 return None;
3617 }
3618 if source.chars().count() > crate::llm::MAX_CODE_LEN {
3619 return None;
3620 }
3621 let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
3625 let mut data = serde_json::Map::new();
3626 data.insert("tool".into(), Value::from(target.opaque()));
3627 data.insert("source".into(), Value::from(source.clone()));
3628 args.insert("tool".into(), Value::from(target.opaque()));
3629 args.insert("bytes".into(), Value::from(source.len() as u64));
3630 Some(ResolvedProposal {
3631 action: ActionKind::CodeRevision,
3632 proposal: Proposal::Data { data },
3633 rendered: format!(
3634 "Proposed new source for tool {} ({} bytes), gated by evalset {}",
3635 target.opaque(),
3636 source.len(),
3637 evalset
3638 ),
3639 summary_key: "llm.code_revision",
3640 summary_args: args,
3641 rollbackable: true,
3642 evalset_hash: Some(evalset),
3643 importance: 0.8,
3644 fact_fields: None,
3645 extra_statements: Vec::new(),
3646 replay: None,
3647 })
3648 }
3649 }
3650}
3651
3652#[allow(clippy::too_many_arguments)]
3662fn stamp_llm(
3663 model: &str,
3664 d: &crate::llm::LlmDraft,
3665 target_ref: String,
3666 cited: Vec<String>,
3667 resolved: Option<ResolvedProposal>,
3668 confidence: f64,
3669 now_ms: i64,
3670 scope: &[String],
3671) -> Recommendation {
3672 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3673 let guidance = if d.guidance.trim().is_empty() {
3674 None
3675 } else {
3676 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
3677 };
3678 let replay = resolved.as_ref().and_then(|r| r.replay.clone());
3679 let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
3680 Some(mut r) => {
3681 let content = match &r.fact_fields {
3685 Some(fields) => format!(
3686 "{} {}",
3687 fields.get("relation").and_then(Value::as_str).unwrap_or(""),
3688 fields.get("object").and_then(Value::as_str).unwrap_or("")
3689 ),
3690 None => match &r.proposal {
3691 Proposal::Cal { cal } => cal.clone(),
3692 Proposal::Data { data } => Value::Object(data.clone()).to_string(),
3693 Proposal::Edit { diff, .. } => diff.clone(),
3694 },
3695 };
3696 if let Some(mut fields) = r.fact_fields.take() {
3699 fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
3700 let mut statements = vec![cal::add("fact", &fields)];
3701 statements.extend(r.extra_statements.iter().cloned());
3702 r.proposal = Proposal::Cal { cal: cal::batch(&statements) };
3703 }
3704 let mut args = r.summary_args;
3705 args.insert("text".into(), Value::from(summary_text));
3706 (
3707 r.action,
3708 r.proposal,
3709 Summary::new(r.summary_key, args),
3710 r.rollbackable,
3711 r.importance,
3712 r.evalset_hash,
3713 Some(content),
3714 )
3715 }
3716 None => {
3717 let mut args = serde_json::Map::new();
3718 args.insert("text".into(), Value::from(summary_text));
3719 let mut data = serde_json::Map::new();
3720 data.insert("source".into(), Value::from("llm"));
3721 (
3722 ActionKind::Flag,
3723 Proposal::Data { data },
3724 Summary::new("llm.discover", args),
3725 false,
3726 0.3,
3727 None,
3728 None,
3729 )
3730 }
3731 };
3732 let dedup = match &content {
3736 Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
3737 None => dedup_key("llm", &target_ref, action),
3738 };
3739 Recommendation {
3740 hash: String::new(),
3741 analyzer: "loop.llm/1".to_string(),
3742 params_snapshot: serde_json::Map::new(),
3743 origin: Origin::Llm { model: model.to_string() },
3744 target_ref: target_ref.clone(),
3745 action_kind: action,
3746 dedup_key: dedup,
3747 summary,
3748 severity: Severity::Low,
3749 proposal,
3750 destructive: false,
3751 rollbackable,
3752 evidence: cited,
3753 evidence_query: None,
3754 metric: None,
3755 confidence: confidence.clamp(0.0, 1.0),
3757 importance,
3758 created_at_ms: now_ms,
3759 guidance,
3760 evalset_hash,
3761 near_duplicate_of: Vec::new(),
3762 replay,
3763 scope: crate::recommendation::normalize_scope(scope),
3767 judged_by: None,
3768 llm_confidence: None,
3769 status: RecStatus::Pending,
3770 }
3771}
3772
3773fn skill_instructions(min_steps: u32) -> String {
3783 format!(
3784 " (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
3785\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
3786PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
3787calls that reached its goal, which a later session facing the same situation \
3788should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
3789the tool called and the values that mattered (the field checked, the tag set, the \
3790exact format produced), a one-line description, and 'when_to_use' — the situation \
3791that should trigger it. The skill-name is a short identifier (letters, digits, \
3792_ -). If a saved skill already covers this procedure, use ITS name so it is \
3793patched rather than duplicated. Do not propose a skill for a procedure that \
3794failed, or for one already saved and unchanged. A finding that itself describes \
3795two or more steps the agent should carry out in order ('after listing the \
3796tickets, fetch each, then …') IS a procedure: propose it as a skill or a plan, \
3797never as a lesson — a lesson is one rule, and a procedure written as one is a \
3798procedure nobody can open.",
3799 crate::llm::MAX_SKILL_STEPS
3800 )
3801}
3802
3803fn plan_instructions(min_nodes: u32) -> String {
3807 format!(
3808 " (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
3809\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
3810\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
3811with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
3812but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
3813'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
3814the cited evidence — and what the step does with it; and 'edges' from step to \
3815step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
3816'path != literal', 'path exists' or '!path', where path is dotted names and the \
3817literal is a JSON string, number, true, false or null — no other operators; state \
3818a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
3819earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
3820has branches or a loop; prefer a skill when it is a straight list. If a saved \
3821plan already covers this procedure, use ITS name so it is patched.",
3822 crate::llm::MAX_PLAN_NODES
3823 )
3824}
3825
3826struct PlanFields {
3830 skill: serde_json::Map<String, Value>,
3831 workflow: Option<serde_json::Map<String, Value>>,
3835 name: String,
3836 n_nodes: usize,
3837 n_edges: usize,
3838 existing_skill: Option<String>,
3839 existing_plan: Option<String>,
3840}
3841
3842#[allow(clippy::too_many_arguments)]
3852fn derived_plan_fields<S: SubstrateRead>(
3853 sub: &S,
3854 target: &TargetRef,
3855 description: &str,
3856 when_to_use: &str,
3857 nodes: &[crate::llm::PlanNodeDraft],
3858 edges: &[crate::llm::PlanEdgeDraft],
3859 cited: &[String],
3860 ns_by_hash: &std::collections::BTreeMap<String, String>,
3861 plans: &crate::policy::PlanAuthoring,
3862) -> Option<PlanFields> {
3863 if target.scheme() != "entity" {
3864 return None;
3865 }
3866 let name = sanitize_skill_name(
3867 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3868 )?;
3869 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3870 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3871 if description.is_empty() || when_to_use.is_empty() {
3872 return None;
3873 }
3874 if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
3875 return None;
3876 }
3877 let known_tools: BTreeSet<String> = cited
3879 .iter()
3880 .filter_map(|h| sub.grain(h).ok().flatten())
3881 .filter_map(|g| g.tool_name().map(normalize_ident))
3882 .collect();
3883 let mut ids: Vec<String> = Vec::new();
3884 let mut steps: Vec<String> = Vec::new();
3885 let mut seen: BTreeSet<String> = BTreeSet::new();
3886 let mut grounded = 0usize;
3887 for n in nodes {
3888 let id = sanitize_skill_name(&n.id)?;
3889 if !seen.insert(id.clone()) {
3890 return None; }
3892 let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
3893 if step.is_empty() {
3894 return None;
3895 }
3896 let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
3906 if !tool.is_empty() && known_tools.contains(&normalize_ident(&tool)) {
3907 grounded += 1;
3908 steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
3909 } else {
3910 steps.push(format!("{}. {id}: {step}", steps.len() + 1));
3911 }
3912 ids.push(id);
3913 }
3914 if grounded == 0 {
3917 return None;
3918 }
3919 let mut edge_vals: Vec<Value> = Vec::new();
3924 let mut flow_lines: Vec<String> = Vec::new();
3925 let mut runnable = true;
3926 for e in edges {
3927 let (Some(src), Some(dst)) = (sanitize_skill_name(&e.src), sanitize_skill_name(&e.dst)) else {
3928 runnable = false;
3929 continue;
3930 };
3931 if !seen.contains(&src) || !seen.contains(&dst) {
3932 flow_lines.push(format!("{} → {}", e.src.trim(), e.dst.trim()));
3934 runnable = false;
3935 continue;
3936 }
3937 let mut ev = serde_json::Map::new();
3938 ev.insert("src".into(), Value::from(src.clone()));
3939 ev.insert("dst".into(), Value::from(dst.clone()));
3940 let mut label = format!("{src} → {dst}");
3941 if let Some(c) = e
3942 .cond
3943 .as_deref()
3944 .map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
3945 .filter(|c| !c.is_empty())
3946 {
3947 label.push_str(&format!(" if {c}"));
3948 ev.insert("cond".into(), Value::from(c));
3949 }
3950 if let Some(m) = e.max_cycles {
3951 if m == 0 || m > 100 {
3952 runnable = false;
3953 } else {
3954 label.push_str(&format!(" (at most {m} times)"));
3955 ev.insert("max_cycles".into(), Value::from(m));
3956 }
3957 }
3958 flow_lines.push(label);
3959 edge_vals.push(Value::Object(ev));
3960 }
3961 if edge_vals.len() > 4 * ids.len() {
3962 runnable = false;
3963 }
3964 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3966 for h in cited {
3967 if let Some(ns) = ns_by_hash.get(h) {
3968 if !ns.is_empty() {
3969 *ns_counts.entry(ns.as_str()).or_default() += 1;
3970 }
3971 }
3972 }
3973 let ns = ns_counts
3974 .iter()
3975 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3976 .map(|(ns, _)| ns.to_string());
3977
3978 let mut workflow = serde_json::Map::new();
3981 workflow.insert("nodes".into(), Value::from(ids.clone()));
3982 workflow.insert("edges".into(), Value::Array(edge_vals));
3983 workflow.insert("name".into(), Value::from(name.clone()));
3984 if let Some(ns) = &ns {
3985 workflow.insert("namespace".into(), Value::from(ns.clone()));
3986 }
3987 let workflow = (runnable && sub.validate_plan(&Value::Object(workflow.clone())).is_ok())
3991 .then_some(workflow);
3992
3993 let mut instructions = steps.join("\n");
3996 if !flow_lines.is_empty() {
3997 instructions.push_str("\n\nFlow:\n");
3998 instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
3999 }
4000 let mut skill = serde_json::Map::new();
4001 skill.insert("name".into(), Value::from(name.clone()));
4002 skill.insert("description".into(), Value::from(description));
4003 skill.insert("when_to_use".into(), Value::from(when_to_use));
4004 skill.insert("instructions".into(), Value::from(instructions));
4005 if let Some(ns) = &ns {
4006 skill.insert("namespace".into(), Value::from(ns.clone()));
4007 }
4008 let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
4009 sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
4010 .ok()?
4011 .into_iter()
4012 .find(|g| pick(g))
4013 .map(|g| g.hash)
4014 };
4015 let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
4016 let existing_plan = workflow
4017 .is_some()
4018 .then(|| live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str())))
4019 .flatten();
4020 let n_edges = workflow
4021 .as_ref()
4022 .and_then(|w| w.get("edges"))
4023 .and_then(Value::as_array)
4024 .map_or(0, |a| a.len());
4025 Some(PlanFields { skill, workflow, name, n_nodes: ids.len(), n_edges, existing_skill, existing_plan })
4026}
4027
4028pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
4030
4031fn detect_premise_drift<S: OmsSubstrate>(
4047 sub: &S,
4048 p: &mut LoopPersisted,
4049 now_ms: i64,
4050) -> Result<Vec<OutcomeInput>> {
4051 let mut out = Vec::new();
4052 let applied: Vec<(String, String, Vec<String>)> = p
4053 .applied
4054 .iter()
4055 .filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
4056 .map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
4057 .collect();
4058 for (rec_hash, target_ref, own) in applied {
4059 let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
4060 if rec.evidence.is_empty() {
4061 continue;
4062 }
4063 let mut moved = 0u64;
4064 for e in &rec.evidence {
4065 match sub.grain(e)? {
4066 None => moved += 1, Some(g) => {
4068 let Some(newer) = &g.superseded_by else { continue };
4069 if own.iter().any(|c| c == newer) {
4070 continue; }
4072 match sub.grain(newer)? {
4073 None => moved += 1,
4076 Some(n) => {
4077 if !same_value(&g, &n) {
4078 moved += 1;
4079 }
4080 }
4081 }
4082 }
4083 }
4084 }
4085 if moved == 0 {
4086 continue;
4087 }
4088 let already = p
4089 .outcomes
4090 .get(&rec_hash)
4091 .and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
4092 .is_some_and(|o| o.current == moved as f64);
4093 if !already {
4094 p.outcomes.entry(rec_hash.clone()).or_default().push(
4095 crate::recommendation::OutcomeResult {
4096 rec_hash: rec_hash.clone(),
4097 metric: PREMISE_DRIFT_METRIC.into(),
4098 baseline: 0.0,
4099 current: moved as f64,
4100 verdict: "drifted".into(),
4101 baseline_kind: "snapshot".into(),
4102 baseline_run_id: None,
4103 best_before: None,
4104 tolerance: 0.0,
4105 current_run_id: None,
4106 cost: None,
4107 horizon_ms: 0,
4108 checkpoint: None,
4109 measured_at_ms: now_ms,
4110 },
4111 );
4112 }
4113 out.push(OutcomeInput {
4114 rec_hash,
4115 target_ref,
4116 metric: PREMISE_DRIFT_METRIC.into(),
4117 baseline: 0.0,
4118 current: moved as f64,
4119 unit: "superseded premises".into(),
4120 higher_is_better: false,
4121 baseline_kind: "snapshot".into(),
4122 baseline_run_id: None,
4123 best_before: None,
4124 tolerance: 0.0,
4125 current_run_id: None,
4126 cost: None,
4127 });
4128 }
4129 Ok(out)
4130}
4131
4132fn moved_premises<S: OmsSubstrate>(
4137 sub: &S,
4138 rec: &Recommendation,
4139 own: &[String],
4140) -> Result<(u64, u64)> {
4141 let total = rec.evidence.len() as u64;
4142 let mut moved = 0u64;
4143 for e in &rec.evidence {
4144 match sub.grain(e)? {
4145 None => moved += 1, Some(g) => {
4147 let Some(newer) = &g.superseded_by else { continue };
4148 if own.iter().any(|c| c == newer) {
4149 continue;
4150 }
4151 match sub.grain(newer)? {
4152 None => moved += 1,
4153 Some(n) => {
4154 if !same_value(&g, &n) {
4155 moved += 1;
4156 }
4157 }
4158 }
4159 }
4160 }
4161 }
4162 Ok((moved, total))
4163}
4164
4165fn withdraw_drifted_open<S: OmsSubstrate>(
4183 sub: &mut S,
4184 p: &mut LoopPersisted,
4185 require_all: bool,
4186 now_ms: i64,
4187) -> Result<u64> {
4188 let open: Vec<String> = p
4189 .status_index
4190 .iter()
4191 .filter(|(_, st)| matches!(st, RecStatus::Pending | RecStatus::Approved))
4192 .map(|(h, _)| h.clone())
4193 .collect();
4194 let mut withdrawn = 0u64;
4195 for rec_hash in open {
4196 let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
4197 if rec.evidence.is_empty() {
4198 continue;
4199 }
4200 let (moved, total) = moved_premises(sub, &rec, &[])?;
4201 let enough = if require_all { moved >= total } else { moved > 0 };
4202 if moved == 0 || !enough {
4203 continue;
4204 }
4205 let from = p.status_index.get(&rec_hash).copied().unwrap_or(RecStatus::Pending);
4206 let prev = p.audit_heads.get(&rec_hash).cloned();
4207 let audit = AuditRecord {
4208 rec_hash: rec_hash.clone(),
4209 from: Some(from),
4210 to: RecStatus::Withdrawn,
4211 actor: "engine:loop.premise_drift".into(),
4212 observer_type: ObserverType::System,
4213 because: format!(
4214 "{moved} of {total} cited grains were superseded by a different value or retracted"
4215 ),
4216 previous_audit_hash: prev,
4217 gating: None,
4218 at_ms: now_ms,
4219 };
4220 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
4221 p.audit_heads.insert(rec_hash.clone(), audit_hash);
4222 p.status_index.insert(rec_hash, RecStatus::Withdrawn);
4223 withdrawn += 1;
4224 }
4225 Ok(withdrawn)
4226}
4227
4228fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
4233 if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
4234 return normalize_ident(a) == normalize_ident(b);
4235 }
4236 for key in ["content", "tool_content", "body", "text", "object"] {
4237 if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
4238 return normalize_ident(a) == normalize_ident(b);
4239 }
4240 }
4241 false
4242}
4243
4244fn sanitize_skill_name(s: &str) -> Option<String> {
4246 let t = s.trim();
4247 if t.is_empty()
4248 || t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
4249 || !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
4250 {
4251 return None;
4252 }
4253 Some(t.to_string())
4254}
4255
4256struct SkillFields {
4260 fields: serde_json::Map<String, Value>,
4261 name: String,
4262 n_steps: usize,
4263 existing: Option<String>,
4264}
4265
4266#[allow(clippy::too_many_arguments)]
4271fn derived_skill_fields<S: SubstrateRead>(
4272 sub: &S,
4273 target: &TargetRef,
4274 description: &str,
4275 when_to_use: &str,
4276 steps: &[String],
4277 cited: &[String],
4278 ns_by_hash: &std::collections::BTreeMap<String, String>,
4279 skills: &crate::policy::SkillAuthoring,
4280) -> Option<SkillFields> {
4281 if target.scheme() != "entity" {
4282 return None;
4283 }
4284 let name = sanitize_skill_name(
4285 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
4286 )?;
4287 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
4288 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
4289 let steps: Vec<String> = steps
4290 .iter()
4291 .map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
4292 .filter(|st| !st.is_empty())
4293 .take(crate::llm::MAX_SKILL_STEPS)
4294 .collect();
4295 if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
4296 return None;
4297 }
4298 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
4301 for h in cited {
4302 if let Some(ns) = ns_by_hash.get(h) {
4303 if !ns.is_empty() {
4304 *ns_counts.entry(ns.as_str()).or_default() += 1;
4305 }
4306 }
4307 }
4308 let ns = ns_counts
4309 .iter()
4310 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
4311 .map(|(ns, _)| ns.to_string());
4312 let instructions = steps
4313 .iter()
4314 .enumerate()
4315 .map(|(i, st)| format!("{}. {st}", i + 1))
4316 .collect::<Vec<_>>()
4317 .join("\n");
4318 let mut fields = serde_json::Map::new();
4319 fields.insert("name".into(), Value::from(name.clone()));
4320 fields.insert("description".into(), Value::from(description));
4321 fields.insert("when_to_use".into(), Value::from(when_to_use));
4322 fields.insert("instructions".into(), Value::from(instructions));
4323 if let Some(ns) = &ns {
4324 fields.insert("namespace".into(), Value::from(ns.clone()));
4325 }
4326 let existing = sub
4328 .grains_of_type(
4329 crate::model::grain_type::SKILL,
4330 ns.as_deref(),
4331 ReadOpts { live_only: true, since_ms: None },
4332 )
4333 .ok()?
4334 .into_iter()
4335 .find(|g| g.skill_name() == Some(name.as_str()))
4336 .map(|g| g.hash);
4337 Some(SkillFields { fields, name, n_steps: steps.len(), existing })
4338}
4339
4340fn derived_fact_fields(
4341 target: &TargetRef,
4342 relation: &str,
4343 object: &str,
4344 cited: &[String],
4345 ns_by_hash: &std::collections::BTreeMap<String, String>,
4346) -> Option<serde_json::Map<String, Value>> {
4347 if target.scheme() != "entity" {
4348 return None;
4349 }
4350 let subject = target
4351 .opaque()
4352 .rsplit_once('/')
4353 .map(|(_, s)| s)
4354 .unwrap_or(target.opaque());
4355 if subject.is_empty() {
4356 return None;
4357 }
4358 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
4359 for h in cited {
4360 if let Some(ns) = ns_by_hash.get(h) {
4361 if !ns.is_empty() {
4362 *ns_counts.entry(ns.as_str()).or_default() += 1;
4363 }
4364 }
4365 }
4366 let lesson_ns = ns_counts
4367 .iter()
4368 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
4369 .map(|(ns, _)| ns.to_string());
4370 let mut fields = serde_json::Map::new();
4371 fields.insert("subject".into(), Value::from(subject));
4372 fields.insert("relation".into(), Value::from(relation));
4373 fields.insert("object".into(), Value::from(object));
4374 if let Some(ns) = lesson_ns {
4378 fields.insert("namespace".into(), Value::from(ns));
4379 }
4380 Some(fields)
4381}
4382
4383fn latest_verdicts(p: &LoopPersisted) -> BTreeMap<String, String> {
4388 let mut out = BTreeMap::new();
4389 for (rec_hash, applied) in &p.applied {
4390 let latest = p
4391 .outcomes
4392 .get(rec_hash)
4393 .and_then(|v| v.iter().max_by_key(|o| o.measured_at_ms))
4394 .map(|o| o.verdict.clone());
4395 for h in &applied.created_hashes {
4396 out.insert(h.clone(), latest.clone().unwrap_or_else(|| "unmeasured".into()));
4397 }
4398 }
4399 out
4400}
4401
4402pub const NEAR_DUPLICATE_COSINE: f64 = 0.90;
4405pub const NEAR_DUPLICATE_JACCARD: f64 = 0.60;
4408const NEAR_DUPLICATE_CAP: usize = 8;
4410
4411pub(crate) fn near_duplicates_of<S: SubstrateRead + ?Sized>(
4416 sub: &S,
4417 subject: &str,
4418 namespace: Option<&str>,
4419 text: &str,
4420) -> Vec<crate::recommendation::NearDuplicate> {
4421 use crate::analyzers::duplicate_sweep::{jaccard, tokenize};
4422 if subject.is_empty() || text.trim().is_empty() {
4423 return Vec::new();
4424 }
4425 let Ok(facts) = sub.grains_of_type(
4426 crate::model::grain_type::FACT,
4427 namespace,
4428 ReadOpts { live_only: true, since_ms: None },
4429 ) else {
4430 return Vec::new();
4431 };
4432 let mine = sub.embed(text).ok().flatten();
4433 let my_tokens = tokenize(text);
4434 let mut out: Vec<crate::recommendation::NearDuplicate> = facts
4435 .iter()
4436 .filter(|f| f.fact_relation() == Some("lesson"))
4437 .filter(|f| f.fact_subject().is_some_and(|s| normalize_ident(s) == normalize_ident(subject)))
4438 .filter_map(|f| {
4439 let other = f.fact_object()?;
4440 let (score, method, floor) = match (&mine, sub.embed(other).ok().flatten()) {
4441 (Some(a), Some(b)) => (cosine(a, &b), "cosine", NEAR_DUPLICATE_COSINE),
4442 _ => (jaccard(&my_tokens, &tokenize(other)), "jaccard", NEAR_DUPLICATE_JACCARD),
4443 };
4444 (score >= floor).then(|| crate::recommendation::NearDuplicate {
4445 hash: f.hash.clone(),
4446 score: (score * 1000.0).round() / 1000.0,
4447 method: method.into(),
4448 })
4449 })
4450 .collect();
4451 out.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal).then(a.hash.cmp(&b.hash)));
4452 out.truncate(NEAR_DUPLICATE_CAP);
4453 out
4454}
4455
4456fn cosine(a: &[f32], b: &[f32]) -> f64 {
4457 if a.len() != b.len() || a.is_empty() {
4458 return 0.0;
4459 }
4460 let (mut dot, mut na, mut nb) = (0f64, 0f64, 0f64);
4461 for (x, y) in a.iter().zip(b) {
4462 dot += *x as f64 * *y as f64;
4463 na += *x as f64 * *x as f64;
4464 nb += *y as f64 * *y as f64;
4465 }
4466 if na == 0.0 || nb == 0.0 {
4467 0.0
4468 } else {
4469 dot / (na.sqrt() * nb.sqrt())
4470 }
4471}
4472
4473const CONSOLIDATION_INSTRUCTIONS: &str = " (8) {\"kind\":\"consolidation\",\"lesson\":\"...\",\
4477\"supersedes\":[\"<hash>\",...]} with the same entity target — ONLY in answer to a \
4478'Lesson pile' finding, which lists the live lessons on one entity that exceed \
4479its budget. Write ONE short imperative rule (max 240 chars) that says what \
4480those lessons say together, dropping nothing a lesson that measured 'held' \
4481required and keeping nothing only a lesson that measured 'regressed' or \
4482'drifted' added; 'supersedes' MUST be exactly the hashes that finding lists \
4483(cite them as evidence too). Applying it replaces every listed lesson with \
4484the one line; the reviewer can restore them all.";
4485
4486fn requires_gating(kind: ActionKind) -> bool {
4491 matches!(
4492 kind,
4493 ActionKind::CodeRevision | ActionKind::AdapterRevision
4494 )
4495}
4496
4497fn stamp(
4498 m: &AnalyzerManifest,
4499 params: &crate::manifest::Params,
4500 d: crate::recommendation::RecDraft,
4501 now_ms: i64,
4502 scope: &[String],
4503) -> Result<Recommendation> {
4504 let target = TargetRef::parse(&d.target_ref)?;
4505 crate::recommendation::validate_code_rules(
4509 d.action_kind,
4510 target.target_class(),
4511 d.evalset_hash.as_deref(),
4512 )?;
4513 let revert_of = match (&d.action_kind, &d.proposal) {
4516 (ActionKind::Revert, Proposal::Data { data }) => {
4517 data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
4518 }
4519 _ => None,
4520 };
4521 let dedup = match revert_of.as_deref() {
4522 Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
4523 None => dedup_key(m.family(), &d.target_ref, d.action_kind),
4524 };
4525 let destructive = match &d.proposal {
4526 Proposal::Cal { cal } => cal::contains_destructive(cal),
4527 _ => false,
4528 };
4529 let rollbackable = match &d.proposal {
4530 Proposal::Cal { .. } => !destructive,
4531 Proposal::Edit { .. } => false,
4532 Proposal::Data { .. } => requires_gating(d.action_kind),
4536 };
4537 let mut evidence = d.evidence;
4538 evidence.truncate(MAX_EVIDENCE);
4539 let origin = match m.trust_class {
4544 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
4545 _ => Origin::Builtin,
4546 };
4547 Ok(Recommendation {
4548 hash: String::new(),
4549 analyzer: m.id.clone(),
4550 params_snapshot: params.snapshot(),
4551 origin,
4552 target_ref: target.as_string(),
4553 action_kind: d.action_kind,
4554 dedup_key: dedup,
4555 summary: d.summary,
4556 severity: d.severity,
4557 proposal: d.proposal,
4558 destructive,
4559 rollbackable,
4560 evidence,
4561 evidence_query: d.evidence_query,
4562 metric: d.metric,
4563 confidence: d.confidence,
4564 importance: d.importance,
4565 created_at_ms: now_ms,
4566 guidance: None,
4567 evalset_hash: d.evalset_hash,
4568 near_duplicate_of: Vec::new(),
4569 replay: None,
4570 scope: crate::recommendation::normalize_scope(scope),
4575 judged_by: d.judged_by,
4578 llm_confidence: None,
4579 status: RecStatus::Pending,
4580 })
4581}
4582
4583fn validate_because(because: &str) -> Result<String> {
4584 let trimmed = because.trim();
4585 if trimmed.is_empty() {
4586 return Err(Error::InvalidProposal(
4587 "a BECAUSE reason is required".into(),
4588 ));
4589 }
4590 if trimmed.chars().count() > MAX_BECAUSE {
4591 return Err(Error::InvalidProposal(format!(
4592 "BECAUSE exceeds {MAX_BECAUSE} chars"
4593 )));
4594 }
4595 Ok(trimmed.to_string())
4596}
4597
4598fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
4599 for req in &m.requires {
4600 match req {
4601 Capability::Forks if !caps.forks => return Some("forks"),
4602 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
4603 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
4604 _ => {}
4605 }
4606 }
4607 None
4608}
4609
4610fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
4611 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
4612}
4613
4614fn gate(
4615 opts: &RunOptions,
4616 p: &LoopPersisted,
4617 new_grains: u64,
4618 new_errors: u64,
4619 now_ms: i64,
4620) -> Option<SkipReason> {
4621 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
4622 if !any {
4623 return None;
4624 }
4625 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
4626 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
4627 let stale_ok = opts
4628 .if_stale_ms
4629 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
4630 if min_new_ok || min_err_ok || stale_ok {
4631 return None;
4632 }
4633 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
4635 Some(SkipReason::NotStale)
4636 } else {
4637 Some(SkipReason::MinNewNotMet)
4638 }
4639}
4640
4641#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
4643pub(crate) struct NewSince {
4644 pub grains: u64,
4646 pub error_events: u64,
4648 pub events: u64,
4650 pub sessions: u64,
4652}
4653
4654fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
4655 let opts = ReadOpts {
4656 live_only: false,
4657 since_ms: watermark.map(|w| w + 1),
4658 };
4659 let mut n = NewSince::default();
4660 let mut sessions: BTreeSet<&str> = BTreeSet::new();
4661 let mut events_held: Vec<GrainRecord> = Vec::new();
4662 for t in [
4663 crate::model::grain_type::FACT,
4664 crate::model::grain_type::EVENT,
4665 crate::model::grain_type::TOOL,
4666 crate::model::grain_type::OBSERVATION,
4667 ] {
4668 let g = sub.grains_of_type(t, None, opts)?;
4669 n.grains += g.len() as u64;
4670 if t == crate::model::grain_type::TOOL {
4672 n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
4673 }
4674 if t == crate::model::grain_type::EVENT {
4675 n.events = g.len() as u64;
4676 events_held = g;
4677 }
4678 }
4679 for e in &events_held {
4680 if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
4681 sessions.insert(sid);
4682 }
4683 }
4684 n.sessions = sessions.len() as u64;
4685 Ok(n)
4686}
4687
4688fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
4695 if !c.is_set() {
4696 return None;
4697 }
4698 let time_ok = c
4699 .every_ms
4700 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
4701 let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
4702 let events_ok = c.every_events.is_some_and(|m| new.events >= m);
4703 let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
4704 if time_ok || grains_ok || events_ok || sessions_ok {
4705 None
4706 } else {
4707 Some(SkipReason::CadenceNotDue)
4708 }
4709}
4710
4711fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
4712 let grains = sub.grains_of_type(
4713 crate::model::grain_type::RECOMMENDATION,
4714 Some(LOOP_NS),
4715 ReadOpts {
4716 live_only: false,
4717 since_ms: None,
4718 },
4719 )?;
4720 let mut set = BTreeSet::new();
4721 for g in grains {
4722 let status = p
4723 .status_index
4724 .get(&g.hash)
4725 .copied()
4726 .unwrap_or(RecStatus::Pending);
4727 if matches!(
4741 status,
4742 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
4743 ) {
4744 if let Some(key) = g.str_field("dedup_key") {
4745 set.insert(key.to_string());
4746 }
4747 }
4748 }
4749 Ok(set)
4750}
4751
4752pub(crate) fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
4759 const BASE_MS: i64 = 7 * 86_400_000;
4760 const CAP_MS: i64 = 90 * 86_400_000;
4761 let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
4762 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
4763 *strikes = strikes.saturating_add(1);
4764 p.cooldowns.insert(dedup_key, now_ms + interval);
4765}
4766
4767fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
4768 let g = sub
4769 .grain(rec_hash)?
4770 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
4771 Recommendation::from_fields(rec_hash, &g.fields)
4772}
4773
4774pub(crate) fn is_definition_statement(line: &str) -> bool {
4782 let up = line.trim_start().to_ascii_uppercase();
4783 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
4784}
4785
4786const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
4789 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
4790 approve it to acknowledge it and let it expire.";
4791
4792const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
4797 (evalset hash + run id + stats) — use apply_gated";
4798
4799const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
4800 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
4801 acknowledge it and let it expire.";
4802
4803pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
4816 match proposal {
4817 Proposal::Cal { .. } => Ok(()),
4818 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
4819 Proposal::Data { data } => {
4825 if requires_gating(action_kind)
4826 || data.get("revert_of").and_then(Value::as_str).is_some()
4827 {
4828 Ok(())
4829 } else {
4830 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
4831 }
4832 }
4833 }
4834}
4835
4836#[cfg(test)]
4837mod definition_body_tests {
4838 use super::safe_definition_body;
4839
4840 #[test]
4841 fn ordinary_bodies_pass() {
4842 assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
4843 assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
4844 }
4845
4846 #[test]
4847 fn a_body_cannot_close_its_own_block_or_carry_destruction() {
4848 assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
4853 assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
4854 assert!(!safe_definition_body("RECALL facts FORGET abc"));
4856 assert!(!safe_definition_body("recall facts purge older than 1d"));
4857 assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
4858 assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
4860 assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
4862 }
4863}
4864
4865#[cfg(test)]
4866mod plan_edit_tests {
4867 use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
4868 use serde_json::json;
4869
4870 fn plan() -> serde_json::Value {
4871 json!({
4872 "nodes": ["fetch", "review", "post"],
4873 "edges": [
4874 {"src": "fetch", "dst": "review"},
4875 {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
4876 ],
4877 "bindings": {"fetch": "sha256:tool1"},
4878 "retries": {"fetch": 1}
4879 })
4880 }
4881
4882 #[test]
4883 fn the_allowlist_admits_thresholds_and_refuses_topology() {
4884 assert!(plan_edit_allowed("edges.1.cond"));
4885 assert!(plan_edit_allowed("edges.1.max_cycles"));
4886 assert!(plan_edit_allowed("retries.fetch"));
4887 for path in [
4890 "nodes",
4891 "nodes.0",
4892 "edges.0.src",
4893 "edges.0.dst",
4894 "edges",
4895 "bindings.fetch",
4896 "edges.x.cond",
4897 "",
4898 ] {
4899 assert!(!plan_edit_allowed(path), "{path} must not be editable");
4900 }
4901 }
4902
4903 #[test]
4904 fn values_are_type_checked_against_the_field() {
4905 assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
4908 assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
4909 assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
4910 assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
4911 assert!(plan_value_ok("retries.fetch", &json!(3)));
4912 assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
4913 assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
4914 assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
4915 assert!(!plan_value_ok("edges.0.src", &json!("other")));
4916 }
4917
4918 #[test]
4919 fn get_reads_through_arrays_and_objects_and_absence_is_null() {
4920 let p = plan();
4921 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
4922 assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
4923 assert_eq!(plan_get(&p, "retries.review"), json!(null));
4926 assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
4927 assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
4928 }
4929
4930 #[test]
4931 fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
4932 let mut p = plan();
4933 assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
4934 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
4935 assert!(plan_set(&mut p, "retries.review", json!(2)));
4936 assert_eq!(plan_get(&p, "retries.review"), json!(2));
4937 assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
4938 assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
4939 let mut q = plan();
4944 q.as_object_mut().unwrap().remove("retries");
4945 assert!(plan_set(&mut q, "retries.greet", json!(1)));
4946 assert_eq!(plan_get(&q, "retries.greet"), json!(1));
4947 let mut r = plan();
4949 r.as_object_mut().unwrap().remove("edges");
4950 assert!(!plan_set(&mut r, "edges.0.cond", json!("x")));
4951 }
4952}
4953
4954#[cfg(test)]
4955mod definition_proposal_tests {
4956 use super::is_definition_statement;
4957
4958 #[test]
4959 fn definition_statements_are_recognized_in_both_spellings() {
4960 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
4961 assert!(is_definition_statement(" define template foo AS { x }"));
4962 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
4963 assert!(!is_definition_statement("ADD fact {}"));
4965 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
4966 assert!(!is_definition_statement("FORGET abc"));
4967 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
4971 }
4972}