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}
141
142#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
153pub struct LlmFunnel {
154 pub evidence: u64,
156 pub proposed: u64,
158 pub cited: u64,
160 pub dropped_uncited: u64,
164 pub dropped_target: u64,
167 pub grounded: u64,
169 pub ground_verdicts: u64,
174 pub ground_call_failed: bool,
178 pub kept: u64,
180 pub stored: u64,
182 #[serde(default, skip_serializing_if = "is_zero")]
186 pub advisory_thin_evidence: u64,
187 #[serde(default, skip_serializing_if = "is_zero")]
193 pub dropped_near_duplicate: u64,
194}
195
196fn is_zero(n: &u64) -> bool {
197 *n == 0
198}
199
200impl RunResult {
201 fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
202 RunResult {
203 outcome: RunOutcome::Skipped,
204 skip_reason: Some(reason),
205 new_grains,
206 new_error_events,
207 proposed: 0,
208 deduped: 0,
209 stored: 0,
210 auto_applied: 0,
211 llm_funnel: None,
212 analyzers_run: vec![],
213 analyzers_skipped: vec![],
214 withdrawn: 0,
215 }
216 }
217
218 pub fn ran(&self) -> bool {
219 self.outcome == RunOutcome::Ran
220 }
221}
222
223pub struct Engine {
226 analyzers: Vec<Box<dyn Analyzer>>,
227 policy: crate::policy::Policy,
228 llm: Option<Box<dyn crate::llm::LlmBackend>>,
231 ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
236}
237
238pub(crate) struct AnalysisPass {
239 pub(crate) survivors: Vec<Recommendation>,
240 proposed: u64,
241 deduped: u64,
242 analyzers_run: Vec<String>,
243 pub(crate) analyzers_skipped: Vec<AnalyzerSkip>,
244 llm_funnel: Option<LlmFunnel>,
245}
246
247impl Engine {
248 pub fn with_builtins() -> Self {
251 Engine {
252 analyzers: crate::analyzer::builtin_analyzers(),
253 policy: crate::policy::Policy::default(),
254 llm: None,
255 ground_llm: None,
256 }
257 }
258
259 pub fn empty() -> Self {
261 Engine {
262 analyzers: vec![],
263 policy: crate::policy::Policy::default(),
264 llm: None,
265 ground_llm: None,
266 }
267 }
268
269 pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
271 self.policy = policy;
272 self
273 }
274
275 pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
280 self.llm = Some(backend);
281 self
282 }
283
284 pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
288 self.ground_llm = Some(backend);
289 self
290 }
291
292 pub fn policy(&self) -> &crate::policy::Policy {
293 &self.policy
294 }
295
296 pub fn has_llm(&self) -> bool {
298 self.llm.is_some()
299 }
300
301 pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
303 self.analyzers.push(analyzer);
304 }
305
306 pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
307 &self.analyzers
308 }
309
310 pub fn analyze_only<S: OmsSubstrate>(
319 &self,
320 sub: &S,
321 opts: &RunOptions,
322 overrides: &BTreeMap<String, Map<String, Value>>,
323 now_ms: i64,
324 ) -> Result<Vec<Recommendation>> {
325 let persisted = LoopPersisted::from_value(sub.load_state()?)?;
326 let analysis_watermark = if opts.full_sweep {
327 None
328 } else {
329 persisted.state.watermark_ms
330 };
331 Ok(self
332 .analysis_pass(
333 sub,
334 &persisted,
335 opts,
336 overrides,
337 analysis_watermark,
338 now_ms,
339 &[],
340 )?
341 .survivors)
342 }
343
344 pub fn run<S: OmsSubstrate>(
347 &self,
348 sub: &mut S,
349 opts: &RunOptions,
350 now_ms: i64,
351 ) -> Result<RunResult> {
352 let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
353 let watermark = persisted.state.watermark_ms;
354 let analysis_watermark = if opts.full_sweep { None } else { watermark };
359
360 let new = count_new(sub, watermark)?;
361 let (new_grains, new_error_events) = (new.grains, new.error_events);
362 if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
363 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
364 }
365 let flags_set = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
369 if !flags_set && !opts.full_sweep {
370 if let Some(reason) = cadence_gate(&self.policy.cadence, &persisted, new, now_ms) {
371 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
372 }
373 }
374
375 let mut outcome_inputs = measure_outcomes(sub, &mut persisted, &self.policy, now_ms)?;
380 let mut withdrawn = 0u64;
381 if self.policy.premise_drift {
382 outcome_inputs.extend(detect_premise_drift(sub, &mut persisted, now_ms)?);
383 withdrawn = withdraw_drifted_open(
386 sub,
387 &mut persisted,
388 self.policy.premise_drift_open_all,
389 now_ms,
390 )?;
391 }
392
393 let AnalysisPass {
394 survivors,
395 proposed,
396 deduped,
397 analyzers_run,
398 analyzers_skipped,
399 llm_funnel,
400 } = self.analysis_pass(
401 &*sub,
402 &persisted,
403 opts,
404 &BTreeMap::new(),
405 analysis_watermark,
406 now_ms,
407 &outcome_inputs,
408 )?;
409
410 let mut stored = 0u64;
413 let mut auto_applied = 0u64;
414 for mut rec in survivors {
415 let spec = rec.to_grain_spec(LOOP_NS)?;
416 let hash = sub.put_grain(&spec)?;
417 rec.hash = hash.clone();
418 let actor = format!("engine:{}", rec.analyzer);
419 let audit = AuditRecord {
420 rec_hash: hash.clone(),
421 from: None,
422 to: RecStatus::Pending,
423 actor: actor.clone(),
424 observer_type: ObserverType::System,
425 because: "analyzer proposed".into(),
426 previous_audit_hash: None,
427 gating: None,
428 at_ms: now_ms,
429 };
430 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
431 persisted
432 .status_index
433 .insert(hash.clone(), RecStatus::Pending);
434 persisted.creators.insert(hash.clone(), actor);
435 if !matches!(rec.origin, Origin::Builtin) {
440 if let Some(trigger) = &opts.triggering_actor {
441 persisted.co_creators.insert(hash.clone(), trigger.clone());
442 }
443 }
444 persisted.audit_heads.insert(hash.clone(), audit_hash);
445 stored += 1;
446
447 if self.can_auto_apply(&*sub, &rec) {
448 self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
449 auto_applied += 1;
450 }
451 }
452
453 persisted.state.last_run_ms = Some(now_ms);
454 persisted.state.watermark_ms = Some(now_ms);
455 sub.store_state(&persisted.to_value()?)?;
456
457 Ok(RunResult {
458 outcome: RunOutcome::Ran,
459 skip_reason: None,
460 new_grains,
461 new_error_events,
462 proposed,
463 deduped,
464 stored,
465 auto_applied,
466 analyzers_run,
467 analyzers_skipped,
468 llm_funnel,
469 withdrawn,
470 })
471 }
472
473 #[allow(clippy::too_many_arguments)]
477 #[allow(clippy::too_many_arguments)]
478 fn analysis_pass<S: OmsSubstrate>(
479 &self,
480 sub: &S,
481 persisted: &LoopPersisted,
482 opts: &RunOptions,
483 external_overrides: &BTreeMap<String, Map<String, Value>>,
484 analysis_watermark: Option<i64>,
485 now_ms: i64,
486 outcome_inputs: &[OutcomeInput],
487 ) -> Result<AnalysisPass> {
488 let existing = existing_dedup_keys(sub, persisted)?;
489 self.analysis_pass_inner(
490 sub,
491 persisted,
492 &self.policy,
493 opts,
494 external_overrides,
495 analysis_watermark,
496 now_ms,
497 outcome_inputs,
498 &existing,
499 None,
500 )
501 }
502
503 #[allow(clippy::too_many_arguments)]
510 pub(crate) fn analysis_pass_inner<S: OmsSubstrate>(
511 &self,
512 sub: &S,
513 persisted: &LoopPersisted,
514 policy: &crate::policy::Policy,
515 opts: &RunOptions,
516 external_overrides: &BTreeMap<String, Map<String, Value>>,
517 analysis_watermark: Option<i64>,
518 now_ms: i64,
519 outcome_inputs: &[OutcomeInput],
520 existing: &BTreeSet<String>,
521 replay: Option<&str>,
522 ) -> Result<AnalysisPass> {
523 let mut analyzers_run = Vec::new();
524 let mut analyzers_skipped = Vec::new();
525 let mut candidates: Vec<Recommendation> = Vec::new();
526 let caps = sub.capabilities();
527 let verdicts = latest_verdicts(persisted);
528
529 for analyzer in &self.analyzers {
530 let m = analyzer.manifest();
531 if let Some(why) = replay {
532 let out_of_process = m.trust_class == crate::manifest::TrustClass::Command;
533 let telemetry_fed = m.requires.contains(&crate::manifest::Capability::Telemetry);
534 if out_of_process || telemetry_fed {
535 analyzers_skipped.push(AnalyzerSkip {
536 id: m.id.clone(),
537 reason: format!(
538 "not replayed: {}",
539 if out_of_process { why } else { "telemetry rollups are not time-indexed" }
540 ),
541 });
542 continue;
543 }
544 }
545 let cfg = persisted.config.get(&m.id);
546 let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
547 if !enabled {
548 analyzers_skipped.push(AnalyzerSkip {
549 id: m.id.clone(),
550 reason: "disabled".into(),
551 });
552 continue;
553 }
554 if policy.denies(m.family()) {
555 analyzers_skipped.push(AnalyzerSkip {
556 id: m.id.clone(),
557 reason: "denied by host policy".into(),
558 });
559 continue;
560 }
561 if let Some(missing) = missing_capability(m, caps) {
562 analyzers_skipped.push(AnalyzerSkip {
563 id: m.id.clone(),
564 reason: format!("missing capability: {missing}"),
565 });
566 continue;
567 }
568 let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
569 if let Some(extra) = external_overrides.get(&m.id) {
570 for (key, value) in extra {
571 param_overrides.insert(key.clone(), value.clone());
572 }
573 }
574 let params = match m.resolve_params(¶m_overrides) {
575 Ok(p) => p,
576 Err(e) => {
577 analyzers_skipped.push(AnalyzerSkip {
578 id: m.id.clone(),
579 reason: e.to_string(),
580 });
581 continue;
582 }
583 };
584 let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
585 let ns_slice: &[String] = if ns_owned.is_empty() {
586 &opts.namespaces
587 } else {
588 &ns_owned
589 };
590 let reader: &dyn SubstrateRead = sub;
591 let ctx = AnalyzeCtx::new(
592 reader,
593 ¶ms,
594 ns_slice,
595 analysis_watermark,
596 now_ms,
597 outcome_inputs,
598 &verdicts,
599 );
600 match analyzer.analyze(&ctx) {
601 Ok(drafts) => {
602 analyzers_run.push(m.id.clone());
603 for draft in drafts {
604 match stamp(m, ¶ms, draft, now_ms, ns_slice) {
605 Ok(rec) => candidates.push(rec),
606 Err(e) => analyzers_skipped.push(AnalyzerSkip {
607 id: m.id.clone(),
608 reason: e.to_string(),
609 }),
610 }
611 }
612 }
613 Err(e) => analyzers_skipped.push(AnalyzerSkip {
614 id: m.id.clone(),
615 reason: e.to_string(),
616 }),
617 }
618 }
619
620 let mut funnel = LlmFunnel::default();
621 if self.llm.is_some() && replay.is_none() {
622 candidates.extend(self.discover(
623 sub,
624 &candidates,
625 analysis_watermark,
626 &opts.namespaces,
627 now_ms,
628 &mut funnel,
629 ));
630 }
631
632 let proposed = candidates.len() as u64;
633 let mut seen = BTreeSet::new();
634 let mut survivors = Vec::new();
635 for candidate in candidates {
636 let family = crate::manifest::analyzer_family(&candidate.analyzer);
637 let floor = [
638 severity_floor_for(persisted, &candidate.analyzer),
639 policy.severity_floor(family),
640 ]
641 .into_iter()
642 .flatten()
643 .max();
644 if floor.is_some_and(|floor| candidate.severity < floor) {
645 continue;
646 }
647 if !seen.insert(candidate.dedup_key.clone()) {
648 continue;
649 }
650 if existing.contains(&candidate.dedup_key) {
651 continue;
652 }
653 if persisted
654 .cooldowns
655 .get(&candidate.dedup_key)
656 .is_some_and(|until| now_ms < *until)
657 {
658 continue;
659 }
660 survivors.push(candidate);
661 }
662 let deduped = proposed - survivors.len() as u64;
663 if self.llm.is_some() && replay.is_none() {
664 self.enrich(&mut survivors);
665 }
666 Ok(AnalysisPass {
667 survivors,
668 proposed,
669 deduped,
670 analyzers_run,
671 analyzers_skipped,
672 llm_funnel: self.llm.is_some().then_some(funnel),
673 })
674 }
675
676 fn discover<S: OmsSubstrate>(
683 &self,
684 sub: &S,
685 candidates: &[Recommendation],
686 watermark: Option<i64>,
687 namespaces: &[String],
688 now_ms: i64,
689 funnel: &mut LlmFunnel,
690 ) -> Vec<Recommendation> {
691 let Some(llm) = &self.llm else {
692 return Vec::new();
693 };
694 let findings: Vec<crate::llm::FindingBrief> = candidates
695 .iter()
696 .take(32)
697 .map(|c| crate::llm::FindingBrief {
698 analyzer: c.analyzer.clone(),
699 summary: c.summary.render(),
700 target: c.target_ref.clone(),
701 severity: c.severity.as_str().to_string(),
702 })
703 .collect();
704 let attribution = self.policy.evidence_attribution;
710 let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
711 let mut bundle: BTreeSet<String> = BTreeSet::new();
712 let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
713 'cited: for c in candidates {
714 for h in &c.evidence {
715 if evidence.len() >= CITED_SEED_CAP {
724 break 'cited;
725 }
726 if !bundle.contains(h) {
727 if let Ok(Some(g)) = sub.grain(h) {
728 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
729 }
730 }
731 }
732 }
733 let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
734 vec![None]
735 } else {
736 namespaces.iter().map(|n| Some(n.as_str())).collect()
737 };
738 let opts = ReadOpts { live_only: true, since_ms: watermark };
739 let mut tool_seeded = 0usize;
760 let seed_successes = self.policy.skills.enabled || self.policy.plans.enabled;
761 'tools: for want_error in [true, false] {
762 if !want_error && !seed_successes {
763 break;
764 }
765 for ns in &scan_ns {
766 if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
767 for g in recent {
768 if tool_seeded >= TOOL_SEED_CAP
769 || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
770 {
771 break 'tools;
772 }
773 if g.is_error() != want_error {
774 continue;
775 }
776 let before = evidence.len();
777 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
778 if evidence.len() > before {
779 tool_seeded += 1;
780 }
781 }
782 }
783 }
784 }
785 if let Ok(rows) =
794 sub.grains_of_type(crate::model::grain_type::OBSERVATION, Some(crate::eval::HARNESS_NS), opts)
795 {
796 let mut seeded = 0usize;
797 for g in rows {
798 if seeded >= HARNESS_SEED_CAP || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE {
799 break;
800 }
801 let kind = g.str_field("observation_kind").unwrap_or_default();
802 if !HARNESS_EVIDENCE_KINDS.contains(&kind) {
803 continue;
804 }
805 let before = evidence.len();
806 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
807 if evidence.len() > before {
808 seeded += 1;
809 }
810 }
811 }
812 'notes: for ns in &scan_ns {
822 if let Ok(recent) =
823 sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
824 {
825 for g in recent {
826 if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
827 break 'notes;
828 }
829 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
830 }
831 }
832 }
833 'seed: for gt in [
834 crate::model::grain_type::FACT,
835 crate::model::grain_type::OBSERVATION,
836 ] {
837 for ns in &scan_ns {
838 if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
839 for g in recent {
840 if evidence.len() >= EVIDENCE_CAP {
841 break 'seed;
842 }
843 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
844 }
845 }
846 }
847 }
848 funnel.evidence = evidence.len() as u64;
849 if evidence.is_empty() {
850 return Vec::new(); }
852 let (approved, rejected) = self.llm_history(sub);
857 let base = match self.policy.discover_objective {
858 crate::policy::DiscoverObjective::ReviewQueue => DISCOVER_INSTRUCTIONS,
859 crate::policy::DiscoverObjective::Learner => DISCOVER_LEARNER_INSTRUCTIONS,
860 };
861 let mut instructions = base.to_string();
864 if self.policy.skills.enabled {
865 instructions.push_str(&skill_instructions(self.policy.skills.min_steps));
866 }
867 if self.policy.plans.enabled {
868 instructions.push_str(&plan_instructions(self.policy.plans.min_nodes));
869 }
870 if findings.iter().any(|f| f.analyzer.starts_with("loop.lesson_pile/")) {
873 instructions.push_str(CONSOLIDATION_INSTRUCTIONS);
874 }
875 let request = crate::llm::LlmRequest {
876 loop_proto: 1,
877 op: "discover",
878 instructions: &instructions,
879 findings: findings.clone(),
880 evidence: evidence.clone(),
881 rejected,
882 approved,
883 };
884 let Ok(body) = serde_json::to_string(&request) else {
885 return Vec::new();
886 };
887 let raw = match llm.complete(&body) {
888 Ok(r) => r,
889 Err(_) => return Vec::new(), };
891 let caps = sub.capabilities();
895 let mut validated: Vec<ValidatedDraft> = Vec::new();
896 let drafts: Vec<_> = crate::llm::parse_discover(&raw)
897 .recommendations
898 .into_iter()
899 .take(crate::llm::MAX_LLM_DRAFTS)
900 .collect();
901 funnel.proposed = drafts.len() as u64;
902 let id_to_hash: std::collections::BTreeMap<&str, &str> = evidence
903 .iter()
904 .map(|e| (e.id.as_str(), e.hash.as_str()))
905 .collect();
906 for d in drafts {
907 let mut cited: Vec<String> = Vec::new();
908 for c in &d.evidence {
909 if let Some(h) = resolve_citation(c, &bundle, &id_to_hash) {
910 if !cited.contains(&h) {
911 cited.push(h);
912 }
913 }
914 }
915 if cited.is_empty() {
916 funnel.dropped_uncited += 1;
917 continue; }
919 let Ok(target) = TargetRef::parse(&d.target) else {
920 funnel.dropped_target += 1;
921 continue;
922 };
923 let tc = target.target_class();
924 if !matches!(tc, "memory" | "query" | "code") {
930 funnel.dropped_target += 1;
931 continue;
932 }
933 let thin = cited.len() < self.policy.min_evidence as usize;
940 if thin {
941 funnel.advisory_thin_evidence += 1;
942 }
943 let resolved = if thin {
944 None
945 } else {
946 resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps, &self.policy)
947 };
948 if tc == "code" && resolved.is_none() {
953 funnel.dropped_target += 1;
954 continue;
955 }
956 validated.push(ValidatedDraft {
957 draft: d,
958 target_ref: target.as_string(),
959 cited,
960 resolved,
961 });
962 }
963 funnel.cited = validated.len() as u64;
964 if validated.is_empty() {
965 return Vec::new();
966 }
967 let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
974 let outcome_metric = self.outcome_metric_template(sub);
975 self.verify_drafts(
976 sub, &**llm, ground, validated, &evidence, outcome_metric, now_ms, funnel, namespaces,
977 )
978 }
979
980 fn outcome_metric_template<S: OmsSubstrate>(
988 &self,
989 sub: &S,
990 ) -> Option<crate::recommendation::MetricSnapshot> {
991 let e = self.policy.outcome_evalset.as_ref()?;
992 let run = crate::eval::newest_eval_run(sub, &e.hash, None).ok().flatten()?;
993 let baseline = crate::eval::run_value(&run, &e.field)?;
994 let schedule = e.schedule();
998 let ms_only: Vec<i64> = schedule.iter().filter_map(Checkpoint::as_ms).collect();
999 let all_ms = ms_only.len() == schedule.len();
1000 let horizons = if ms_only.is_empty() { vec![86_400_000] } else { ms_only };
1001 Some(crate::recommendation::MetricSnapshot {
1002 metric: format!("evalset:{}:{}", e.hash, e.field),
1003 baseline,
1004 unit: e.field.clone(),
1005 n: run.total(),
1006 window: "per-run".into(),
1007 subject: None,
1008 namespace: None,
1009 relation: None,
1010 query: format!(
1011 "RECALL facts WHERE subject = \"evalset:{}\" AND relation = \"mg:eval_run\"",
1012 e.hash
1013 ),
1014 review_after_ms: horizons[0],
1015 horizons_ms: if all_ms { horizons } else { Vec::new() },
1016 checkpoints: if all_ms { Vec::new() } else { schedule },
1017 higher_is_better: e.higher_is_better,
1018 })
1019 }
1020
1021 #[allow(clippy::too_many_arguments)]
1028 #[allow(clippy::too_many_arguments)]
1029 fn verify_drafts<S: SubstrateRead>(
1030 &self,
1031 sub: &S,
1032 llm: &dyn crate::llm::LlmBackend,
1033 ground: &dyn crate::llm::LlmBackend,
1034 validated: Vec<ValidatedDraft>,
1035 evidence: &[crate::llm::EvidenceItem],
1036 outcome_metric: Option<crate::recommendation::MetricSnapshot>,
1037 now_ms: i64,
1038 funnel: &mut LlmFunnel,
1039 scope: &[String],
1042 ) -> Vec<Recommendation> {
1043 use crate::llm::*;
1044 let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
1045 evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
1046 let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
1047 cited
1048 .iter()
1049 .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
1050 .collect()
1051 };
1052
1053 let claims: Vec<GroundItem> = validated
1057 .iter()
1058 .enumerate()
1059 .map(|(i, v)| GroundItem {
1060 id: i,
1061 claim: claim_text(&v.draft, v.resolved.as_ref()),
1062 evidence: ev_for(&v.cited),
1063 })
1064 .collect();
1065 let ground_req = GroundRequest {
1066 loop_proto: 1,
1067 op: "ground",
1068 instructions: GROUND_INSTRUCTIONS,
1069 claims,
1070 };
1071 let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
1077 .ok()
1078 .and_then(|b| ground.complete(&b).ok())
1079 {
1080 Some(raw) => {
1081 let parsed = parse_ground(&raw);
1082 funnel.ground_verdicts = parsed.results.len() as u64;
1083 parsed
1084 .results
1085 .into_iter()
1086 .filter(|r| r.supported)
1087 .map(|r| r.id)
1088 .collect()
1089 }
1090 None => {
1091 funnel.ground_call_failed = true;
1092 return Vec::new();
1093 }
1094 };
1095 funnel.grounded = grounded.len() as u64;
1096 if grounded.is_empty() {
1097 return Vec::new();
1098 }
1099
1100 let items: Vec<VerifyItem> = validated
1106 .iter()
1107 .enumerate()
1108 .filter(|(i, _)| grounded.contains(i))
1109 .map(|(i, v)| VerifyItem {
1110 id: i,
1111 summary: claim_text(&v.draft, v.resolved.as_ref()),
1113 target: v.target_ref.clone(),
1114 evidence: ev_for(&v.cited),
1115 })
1116 .collect();
1117 let verify_req = VerifyRequest {
1118 loop_proto: 1,
1119 op: "verify",
1120 instructions: VERIFY_INSTRUCTIONS,
1121 findings: items,
1122 };
1123 let verdicts: std::collections::BTreeMap<usize, f64> =
1124 match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
1125 Some(raw) => parse_verify(&raw)
1126 .results
1127 .into_iter()
1128 .filter(|r| r.keep)
1129 .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
1130 .collect(),
1131 None => return Vec::new(),
1132 };
1133
1134 funnel.kept = verdicts.len() as u64;
1135 let mut out = Vec::new();
1139 for (i, v) in validated.into_iter().enumerate() {
1140 if let Some(&conf) = verdicts.get(&i) {
1141 if conf >= MIN_LLM_CONFIDENCE {
1142 let near = match v.resolved.as_ref() {
1151 Some(r) if r.action != ActionKind::Consolidate => r
1152 .fact_fields
1153 .as_ref()
1154 .filter(|f| f.get("relation").and_then(Value::as_str) == Some("lesson"))
1155 .map(|f| {
1156 near_duplicates_of(
1157 sub,
1158 f.get("subject").and_then(Value::as_str).unwrap_or(""),
1159 f.get("namespace").and_then(Value::as_str),
1160 f.get("object").and_then(Value::as_str).unwrap_or(""),
1161 )
1162 })
1163 .unwrap_or_default(),
1164 _ => Vec::new(),
1165 };
1166 if !near.is_empty()
1167 && self.policy.near_duplicate == crate::policy::NearDuplicateMode::Suppress
1168 {
1169 funnel.dropped_near_duplicate += 1;
1170 continue;
1171 }
1172 let mut rec = stamp_llm(
1173 llm.model(),
1174 &v.draft,
1175 v.target_ref,
1176 v.cited,
1177 v.resolved,
1178 conf,
1179 now_ms,
1180 scope,
1181 );
1182 if let Some(best) = near.first() {
1183 rec.summary.args.insert("near_count".into(), Value::from(near.len() as u64));
1185 rec.summary.args.insert("near_score".into(), Value::from(best.score));
1186 rec.summary.args.insert("near_method".into(), Value::from(best.method.clone()));
1187 rec.summary.args.insert(
1188 "near_hash".into(),
1189 Value::from(best.hash.chars().take(12).collect::<String>()),
1190 );
1191 rec.summary.template_id = "llm.lesson_near_duplicate".into();
1192 rec.near_duplicate_of = near;
1193 }
1194 if rec.rollbackable {
1198 rec.metric = outcome_metric.clone();
1199 }
1200 out.push(rec);
1201 }
1202 }
1203 }
1204 funnel.stored = out.len() as u64;
1205 out
1206 }
1207
1208 fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
1213 const MAX: usize = 20;
1214 let Ok(mut recs) = self.recommendations(sub, None) else {
1215 return (Vec::new(), Vec::new());
1216 };
1217 recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
1218 recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
1219 let mut approved = Vec::new();
1220 let mut rejected = Vec::new();
1221 for r in &recs {
1222 match r.status {
1223 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
1224 if approved.len() < MAX =>
1225 {
1226 approved.push(r.summary.render());
1227 }
1228 RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
1229 _ => {}
1230 }
1231 }
1232 (approved, rejected)
1233 }
1234
1235 fn enrich(&self, survivors: &mut [Recommendation]) {
1240 let Some(llm) = &self.llm else {
1241 return;
1242 };
1243 if survivors.is_empty() {
1244 return;
1245 }
1246 let findings: Vec<crate::llm::FindingBrief> = survivors
1247 .iter()
1248 .map(|r| crate::llm::FindingBrief {
1249 analyzer: r.analyzer.clone(),
1250 summary: r.summary.render(),
1251 target: r.target_ref.clone(),
1252 severity: r.severity.as_str().to_string(),
1253 })
1254 .collect();
1255 let request = crate::llm::LlmRequest {
1256 loop_proto: 1,
1257 op: "enrich",
1258 instructions: ENRICH_INSTRUCTIONS,
1259 findings,
1260 evidence: Vec::new(),
1261 rejected: Vec::new(),
1262 approved: Vec::new(),
1263 };
1264 let Ok(body) = serde_json::to_string(&request) else {
1265 return;
1266 };
1267 let raw = match llm.complete(&body) {
1268 Ok(r) => r,
1269 Err(_) => return,
1270 };
1271 for note in crate::llm::parse_enrich(&raw).notes {
1272 if note.guidance.trim().is_empty() {
1273 continue;
1274 }
1275 if let Some(r) = survivors
1276 .iter_mut()
1277 .find(|r| r.target_ref == note.target && r.guidance.is_none())
1278 {
1279 r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
1280 }
1281 }
1282 }
1283
1284 fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
1293 if !rec.origin.auto_apply_eligible() || rec.destructive {
1294 return false;
1295 }
1296 let manifest_ok = self
1300 .analyzers
1301 .iter()
1302 .map(|a| a.manifest())
1303 .find(|m| m.id == rec.analyzer)
1304 .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
1305 if !manifest_ok {
1306 return false;
1307 }
1308 let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1309 return false;
1310 };
1311 let family = crate::manifest::analyzer_family(&rec.analyzer);
1312 if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1313 return false;
1314 }
1315 match &rec.proposal {
1320 Proposal::Cal { cal } => cal
1321 .lines()
1322 .map(str::trim)
1323 .filter(|l| !l.is_empty())
1324 .all(|l| supersede_is_value_identical(sub, l)),
1325 _ => false,
1326 }
1327 }
1328
1329 fn auto_apply<S: OmsSubstrate>(
1332 &self,
1333 sub: &mut S,
1334 p: &mut LoopPersisted,
1335 rec: &Recommendation,
1336 now_ms: i64,
1337 ) -> Result<()> {
1338 let mut created = Vec::new();
1339 if let Proposal::Cal { cal } = &rec.proposal {
1340 if cal.lines().map(str::trim).any(is_definition_statement) {
1345 return Err(Error::InvalidProposal(
1346 "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1347 auto-applied: it changes what every future context contains, so it \
1348 requires a human APPROVE + APPLY with BECAUSE"
1349 .into(),
1350 ));
1351 }
1352 for r in sub.execute_cal(cal)? {
1353 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1354 created.push(h.to_string());
1355 }
1356 }
1357 }
1358 let applied = AppliedRecord {
1359 applied_at_ms: now_ms,
1360 target_ref: rec.target_ref.clone(),
1361 rollbackable: rec.rollbackable,
1362 created_hashes: created,
1363 inverse_cal: None,
1364 metric: rec.metric.clone(),
1365 };
1366 let prev = p.audit_heads.get(&rec.hash).cloned();
1367 let audit = AuditRecord {
1368 rec_hash: rec.hash.clone(),
1369 from: Some(RecStatus::Pending),
1370 to: RecStatus::Applied,
1371 actor: "policy:auto".into(),
1372 observer_type: ObserverType::Policy,
1373 because: "auto-applied per host policy".into(),
1374 previous_audit_hash: prev,
1375 gating: None,
1376 at_ms: now_ms,
1377 };
1378 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1379 p.audit_heads.insert(rec.hash.clone(), audit_hash);
1380 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1381 p.applied.insert(rec.hash.clone(), applied);
1382 Ok(())
1383 }
1384
1385 #[allow(clippy::too_many_arguments)]
1388 pub fn review<S: OmsSubstrate>(
1389 &self,
1390 sub: &mut S,
1391 rec_hash: &str,
1392 decision: Decision,
1393 actor: &str,
1394 observer: ObserverType,
1395 scopes: &ScopeSet,
1396 because: &str,
1397 now_ms: i64,
1398 ) -> Result<()> {
1399 if !scopes.has(Scope::Review) {
1400 return Err(Error::ScopeDenied("review".into()));
1401 }
1402 let because = validate_because(because)?;
1403 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1404 let status = *p
1405 .status_index
1406 .get(rec_hash)
1407 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1408 let to = match decision {
1409 Decision::Approve => RecStatus::Approved,
1410 Decision::Reject => RecStatus::Rejected,
1411 };
1412 if !status.can_transition_to(to, false) {
1413 return Err(Error::LifecycleViolation(format!(
1414 "{} -> {}",
1415 status.as_str(),
1416 to.as_str()
1417 )));
1418 }
1419 if to == RecStatus::Approved {
1420 if let Some(creator) = p.creators.get(rec_hash) {
1421 if creator == actor {
1422 return Err(Error::SelfApproval(format!(
1423 "{actor} created this recommendation"
1424 )));
1425 }
1426 }
1427 if let Some(trigger) = p.co_creators.get(rec_hash) {
1428 if trigger == actor {
1429 return Err(Error::SelfApproval(format!(
1430 "{actor} triggered the run that authored this recommendation"
1431 )));
1432 }
1433 }
1434 }
1435 let prev = p.audit_heads.get(rec_hash).cloned();
1436 let audit = AuditRecord {
1437 rec_hash: rec_hash.into(),
1438 from: Some(status),
1439 to,
1440 actor: actor.into(),
1441 observer_type: observer,
1442 because,
1443 previous_audit_hash: prev,
1444 gating: None,
1445 at_ms: now_ms,
1446 };
1447 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1448 p.audit_heads.insert(rec_hash.into(), audit_hash);
1449 p.status_index.insert(rec_hash.into(), to);
1450 if to == RecStatus::Rejected {
1451 if let Ok(rec) = load_rec(sub, rec_hash) {
1452 strike_cooldown(&mut p, rec.dedup_key, now_ms);
1453 }
1454 }
1455 sub.store_state(&p.to_value()?)?;
1456 Ok(())
1457 }
1458
1459 pub fn preflight_apply<S: OmsSubstrate>(
1476 &self,
1477 sub: &S,
1478 rec_hash: &str,
1479 scopes: &ScopeSet,
1480 allow_destructive: bool,
1481 has_gating: bool,
1482 ) -> Result<()> {
1483 if !scopes.has(Scope::Apply) {
1484 return Err(Error::ScopeDenied("apply".into()));
1485 }
1486 let rec = load_rec(sub, rec_hash)?;
1487 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1488 return Err(Error::DestructiveGated(
1489 "destructive apply requires admin scope + allow_destructive".into(),
1490 ));
1491 }
1492 ensure_executable(rec.action_kind, &rec.proposal)?;
1493 if requires_gating(rec.action_kind) && !has_gating {
1494 return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1495 }
1496 Ok(())
1497 }
1498
1499 #[allow(clippy::too_many_arguments)]
1503 pub fn apply<S: OmsSubstrate>(
1504 &self,
1505 sub: &mut S,
1506 rec_hash: &str,
1507 actor: &str,
1508 observer: ObserverType,
1509 scopes: &ScopeSet,
1510 because: &str,
1511 allow_destructive: bool,
1512 now_ms: i64,
1513 ) -> Result<AppliedRecord> {
1514 self.apply_inner(
1515 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1516 )
1517 }
1518
1519 pub fn gating_evidence<S: OmsSubstrate>(
1526 &self,
1527 sub: &S,
1528 rec_hash: &str,
1529 run_id: &str,
1530 ) -> Result<crate::recommendation::GatingEvidence> {
1531 let rec = self
1532 .recommendations(sub, None)?
1533 .into_iter()
1534 .find(|r| r.hash == rec_hash)
1535 .ok_or_else(|| {
1536 Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1537 })?;
1538 let pin = rec.evalset_hash.ok_or_else(|| {
1539 Error::InvalidProposal(
1540 "this recommendation pins no evalset — a gating run applies only \
1541 to code and adapter revisions"
1542 .into(),
1543 )
1544 })?;
1545 match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1550 Some(run) => Ok(crate::recommendation::GatingEvidence {
1551 evalset_hash: pin,
1552 run_id: run.run_id,
1553 passed: run.passed,
1554 failed: run.failed,
1555 }),
1556 None => Err(Error::InvalidProposal(format!(
1557 "no recorded gate run '{run_id}' for evalset {pin} — run \
1558 `areev eval run --evalset {pin} ...` first"
1559 ))),
1560 }
1561 }
1562
1563 #[allow(clippy::too_many_arguments)]
1567 pub fn apply_gated<S: OmsSubstrate>(
1568 &self,
1569 sub: &mut S,
1570 rec_hash: &str,
1571 actor: &str,
1572 observer: ObserverType,
1573 scopes: &ScopeSet,
1574 because: &str,
1575 allow_destructive: bool,
1576 gating: &crate::recommendation::GatingEvidence,
1577 now_ms: i64,
1578 ) -> Result<AppliedRecord> {
1579 self.apply_inner(
1580 sub,
1581 rec_hash,
1582 actor,
1583 observer,
1584 scopes,
1585 because,
1586 allow_destructive,
1587 Some(gating),
1588 now_ms,
1589 )
1590 }
1591
1592 #[allow(clippy::too_many_arguments)]
1593 fn apply_inner<S: OmsSubstrate>(
1594 &self,
1595 sub: &mut S,
1596 rec_hash: &str,
1597 actor: &str,
1598 observer: ObserverType,
1599 scopes: &ScopeSet,
1600 because: &str,
1601 allow_destructive: bool,
1602 gating: Option<&crate::recommendation::GatingEvidence>,
1603 now_ms: i64,
1604 ) -> Result<AppliedRecord> {
1605 if !scopes.has(Scope::Apply) {
1606 return Err(Error::ScopeDenied("apply".into()));
1607 }
1608 let because = validate_because(because)?;
1609 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1610 let status = *p
1611 .status_index
1612 .get(rec_hash)
1613 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1614 if !status.can_transition_to(RecStatus::Applied, false) {
1615 return Err(Error::LifecycleViolation(format!(
1616 "{} -> applied (approve first)",
1617 status.as_str()
1618 )));
1619 }
1620 let rec = load_rec(sub, rec_hash)?;
1621 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1622 return Err(Error::DestructiveGated(
1623 "destructive apply requires admin scope + allow_destructive".into(),
1624 ));
1625 }
1626 if requires_gating(rec.action_kind) {
1631 let g = gating
1632 .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1633 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1634 if g.evalset_hash != pin {
1635 return Err(Error::InvalidProposal(format!(
1636 "gating ran evalset {} but the recommendation is pinned \
1637 to {pin} (Rule E1)",
1638 g.evalset_hash
1639 )));
1640 }
1641 match sub.grain(pin)? {
1642 Some(evalset) if evalset.is_live() => {}
1643 Some(_) => {
1644 return Err(Error::InvalidProposal(
1645 "the pinned evalset was superseded after gating — \
1646 the recommendation must re-gate (Rule E1)"
1647 .into(),
1648 ))
1649 }
1650 None => {
1651 return Err(Error::InvalidProposal(format!(
1652 "pinned evalset {pin} not found in the substrate"
1653 )))
1654 }
1655 }
1656 if g.failed > 0 {
1657 return Err(Error::InvalidProposal(format!(
1658 "the gating run failed {}/{} cases — a failing gate \
1659 admits nothing",
1660 g.failed,
1661 g.passed + g.failed
1662 )));
1663 }
1664 }
1665
1666 let mut created = Vec::new();
1668 let mut inverse_cal: Option<String> = None;
1672 match &rec.proposal {
1673 Proposal::Cal { cal } => {
1674 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1675 if !is_definition_statement(line) {
1676 continue;
1677 }
1678 match sub.definition_inverse(line)? {
1679 Some(inv) => inverse_cal = Some(inv),
1680 None => {
1681 return Err(Error::InvalidProposal(format!(
1682 "this substrate cannot record a rollback inverse for {line:?}; \
1683 a definition rewrite that ROLLBACK could not undo is refused \
1684 rather than applied"
1685 )))
1686 }
1687 }
1688 }
1689 let rows = sub.execute_cal(cal)?;
1690 for r in rows {
1691 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1692 created.push(h.to_string());
1693 }
1694 }
1695 }
1696 Proposal::Edit { .. } => {
1699 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1700 }
1701 Proposal::Data { data } if requires_gating(rec.action_kind) => {
1709 let relation = if rec.action_kind == ActionKind::AdapterRevision {
1710 "mg:adapter_promotion"
1711 } else {
1712 "mg:code_promotion"
1713 };
1714 let mut promoted = data.clone();
1722 if let Some(Value::String(src)) = promoted.remove("source") {
1723 let address = sub.put_blob(src.as_bytes())?;
1724 promoted.insert("code_address".into(), Value::from(address));
1725 }
1726 let mut spec = crate::substrate::GrainSpec::new(
1727 crate::model::grain_type::FACT,
1728 LOOP_NS,
1729 )
1730 .with_field("subject", rec.target_ref.clone())
1731 .with_field("relation", relation)
1732 .with_field(
1733 "object",
1734 serde_json::to_string(&promoted)
1735 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1736 )
1737 .with_field("rec_hash", rec_hash.to_string());
1738 if let Some(g) = gating {
1739 spec = spec
1740 .with_field("gating_evalset", g.evalset_hash.clone())
1741 .with_field("gating_run_id", g.run_id.clone());
1742 }
1743 created.push(sub.put_grain(&spec)?);
1744 }
1745 Proposal::Data { data } => {
1746 let revert_of = data
1751 .get("revert_of")
1752 .and_then(Value::as_str)
1753 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1754 self.rollback(
1755 sub,
1756 revert_of,
1757 actor,
1758 observer,
1759 scopes,
1760 &because,
1761 now_ms,
1762 )?;
1763 p = LoopPersisted::from_value(sub.load_state()?)?;
1766 if let Ok(reverted) = load_rec(sub, revert_of) {
1776 strike_cooldown(&mut p, reverted.dedup_key, now_ms);
1777 }
1778 }
1779 }
1780
1781 let applied = AppliedRecord {
1782 applied_at_ms: now_ms,
1783 target_ref: rec.target_ref.clone(),
1784 rollbackable: rec.rollbackable,
1785 created_hashes: created,
1786 inverse_cal,
1787 metric: rec.metric.clone(),
1788 };
1789 let prev = p.audit_heads.get(rec_hash).cloned();
1790 let audit = AuditRecord {
1791 rec_hash: rec_hash.into(),
1792 from: Some(status),
1793 to: RecStatus::Applied,
1794 actor: actor.into(),
1795 observer_type: observer,
1796 because,
1797 previous_audit_hash: prev,
1798 gating: gating.cloned(),
1799 at_ms: now_ms,
1800 };
1801 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1802 p.audit_heads.insert(rec_hash.into(), audit_hash);
1803 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1804 p.applied.insert(rec_hash.into(), applied.clone());
1805 sub.store_state(&p.to_value()?)?;
1806 Ok(applied)
1807 }
1808
1809 #[allow(clippy::too_many_arguments)]
1812 pub fn rollback<S: OmsSubstrate>(
1813 &self,
1814 sub: &mut S,
1815 rec_hash: &str,
1816 actor: &str,
1817 observer: ObserverType,
1818 scopes: &ScopeSet,
1819 because: &str,
1820 now_ms: i64,
1821 ) -> Result<()> {
1822 if !scopes.has(Scope::Apply) {
1823 return Err(Error::ScopeDenied("apply".into()));
1824 }
1825 let because = validate_because(because)?;
1826 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1827 let status = *p
1828 .status_index
1829 .get(rec_hash)
1830 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1831 if !status.can_transition_to(RecStatus::RolledBack, false) {
1832 return Err(Error::LifecycleViolation(format!(
1833 "{} -> rolled_back",
1834 status.as_str()
1835 )));
1836 }
1837 let applied = p
1838 .applied
1839 .get(rec_hash)
1840 .cloned()
1841 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1842 if !applied.rollbackable {
1843 return Err(Error::LifecycleViolation(
1844 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1845 ));
1846 }
1847 for h in &applied.created_hashes {
1848 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1849 }
1850 if let Some(inverse) = &applied.inverse_cal {
1857 sub.execute_cal(inverse)?;
1858 }
1859 let prev = p.audit_heads.get(rec_hash).cloned();
1860 let audit = AuditRecord {
1861 rec_hash: rec_hash.into(),
1862 from: Some(status),
1863 to: RecStatus::RolledBack,
1864 actor: actor.into(),
1865 observer_type: observer,
1866 because,
1867 previous_audit_hash: prev,
1868 gating: None,
1869 at_ms: now_ms,
1870 };
1871 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1872 p.audit_heads.insert(rec_hash.into(), audit_hash);
1873 p.status_index
1874 .insert(rec_hash.into(), RecStatus::RolledBack);
1875 sub.store_state(&p.to_value()?)?;
1876 Ok(())
1877 }
1878
1879 pub fn recommendations<S: OmsSubstrate>(
1884 &self,
1885 sub: &S,
1886 status_filter: Option<RecStatus>,
1887 ) -> Result<Vec<Recommendation>> {
1888 let p = LoopPersisted::from_value(sub.load_state()?)?;
1889 let grains = sub.grains_of_type(
1890 crate::model::grain_type::RECOMMENDATION,
1891 Some(LOOP_NS),
1892 ReadOpts {
1893 live_only: false,
1894 since_ms: None,
1895 },
1896 )?;
1897 let mut out = Vec::new();
1898 for g in grains {
1899 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1900 rec.status = p
1901 .status_index
1902 .get(&g.hash)
1903 .copied()
1904 .unwrap_or(RecStatus::Pending);
1905 if let Some(f) = status_filter {
1906 if rec.status != f {
1907 continue;
1908 }
1909 }
1910 out.push(rec);
1911 }
1912 out.sort_by(|a, b| {
1921 b.severity
1922 .cmp(&a.severity)
1923 .then(a.created_at_ms.cmp(&b.created_at_ms))
1924 .then(a.dedup_key.cmp(&b.dedup_key))
1925 .then(a.hash.cmp(&b.hash))
1926 });
1927 Ok(out)
1928 }
1929
1930 pub fn analyzer_settings<S: OmsSubstrate>(
1933 &self,
1934 sub: &S,
1935 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1936 let p = LoopPersisted::from_value(sub.load_state()?)?;
1937 Ok(self
1938 .analyzers
1939 .iter()
1940 .map(|a| {
1941 let m = a.manifest();
1942 let cfg = p.config.get(&m.id);
1943 crate::config::AnalyzerSetting {
1944 id: m.id.clone(),
1945 title: m.title.clone(),
1946 description: m.description.clone(),
1947 tier: format!("{:?}", m.tier),
1948 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1949 default_on: m.default_on,
1950 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1951 severity_floor: cfg
1952 .and_then(|c| c.severity_floor)
1953 .map(|s| s.as_str().to_string()),
1954 }
1955 })
1956 .collect())
1957 }
1958
1959 pub fn set_analyzer_config<S: OmsSubstrate>(
1965 &self,
1966 sub: &mut S,
1967 analyzer_id: &str,
1968 update: crate::config::AnalyzerConfigUpdate,
1969 scopes: &ScopeSet,
1970 ) -> Result<crate::config::AnalyzerConfig> {
1971 if !scopes.has(Scope::Admin) {
1972 return Err(Error::ScopeDenied("admin".into()));
1973 }
1974 let manifest = self
1975 .analyzers
1976 .iter()
1977 .map(|a| a.manifest())
1978 .find(|m| m.id == analyzer_id)
1979 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1980 if let Some(params) = &update.params {
1982 manifest.resolve_params(params)?;
1983 }
1984 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1985 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1986 if let Some(enabled) = update.enabled {
1987 cfg.enabled = Some(enabled);
1988 }
1989 if update.clear_floor {
1990 cfg.severity_floor = None;
1991 } else if let Some(floor) = update.severity_floor {
1992 cfg.severity_floor = Some(floor);
1993 }
1994 if let Some(params) = update.params {
1995 cfg.params = params;
1996 }
1997 if let Some(ns) = update.namespaces {
1998 cfg.namespaces = ns;
1999 }
2000 let stored = cfg.clone();
2001 sub.store_state(&p.to_value()?)?;
2002 Ok(stored)
2003 }
2004
2005 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
2008 let p = LoopPersisted::from_value(sub.load_state()?)?;
2009 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
2010 out.sort_by(|a, b| {
2014 a.measured_at_ms
2015 .cmp(&b.measured_at_ms)
2016 .then(a.horizon_ms.cmp(&b.horizon_ms))
2017 .then(a.metric.cmp(&b.metric))
2018 .then(a.rec_hash.cmp(&b.rec_hash))
2019 });
2020 Ok(out)
2021 }
2022
2023 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
2027 let p = LoopPersisted::from_value(sub.load_state()?)?;
2028 let new = count_new(sub, p.state.watermark_ms)?;
2029 let (grains_since_run, error_events_since_run) = (new.grains, new.error_events);
2030 let recs = self.recommendations(sub, None)?;
2031 let mut pending = 0;
2032 let mut applied = 0;
2033 for r in &recs {
2034 match r.status {
2035 RecStatus::Pending => pending += 1,
2036 RecStatus::Applied => applied += 1,
2037 _ => {}
2038 }
2039 }
2040 let stale = match p.state.last_run_ms {
2042 None => true,
2043 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
2044 };
2045 Ok(Health {
2046 last_run_ms: p.state.last_run_ms,
2047 grains_since_run,
2048 error_events_since_run,
2049 pending,
2050 applied,
2051 total: recs.len() as u64,
2052 stale,
2053 })
2054 }
2055
2056 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
2061 let recs = self.recommendations(sub, None)?;
2062 let mut m = LlmMetrics::default();
2063 for r in &recs {
2064 if !matches!(r.origin, Origin::Llm { .. }) {
2065 continue;
2066 }
2067 m.proposed += 1;
2068 match r.status {
2069 RecStatus::Pending => m.pending += 1,
2070 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
2071 RecStatus::Rejected => m.rejected += 1,
2072 RecStatus::Expired | RecStatus::Withdrawn => {}
2075 }
2076 }
2077 let decided = m.approved + m.rejected;
2078 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
2079 Ok(m)
2080 }
2081}
2082
2083#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
2085pub struct Health {
2086 #[serde(skip_serializing_if = "Option::is_none")]
2087 pub last_run_ms: Option<i64>,
2088 pub grains_since_run: u64,
2089 pub error_events_since_run: u64,
2090 pub pending: u64,
2091 pub applied: u64,
2092 pub total: u64,
2093 pub stale: bool,
2096}
2097
2098#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
2100pub struct LlmMetrics {
2101 pub proposed: u64,
2104 pub pending: u64,
2105 pub approved: u64,
2107 pub rejected: u64,
2108 #[serde(skip_serializing_if = "Option::is_none")]
2110 pub approval_rate: Option<f64>,
2111}
2112
2113fn measure_outcomes<S: OmsSubstrate>(
2120 sub: &S,
2121 p: &mut LoopPersisted,
2122 policy: &crate::policy::Policy,
2123 now_ms: i64,
2124) -> Result<Vec<OutcomeInput>> {
2125 let mut due: Vec<(String, crate::config::AppliedRecord, Checkpoint)> = Vec::new();
2129 for (h, a) in &p.applied {
2130 if p.status_index.get(h) != Some(&RecStatus::Applied) {
2131 continue;
2132 }
2133 let Some(metric) = &a.metric else { continue };
2134 let done = p.measured.get(h).cloned().unwrap_or_default();
2135 for cp in metric.schedule() {
2136 if !done.contains(&cp) && checkpoint_due(sub, metric, a.applied_at_ms, cp, now_ms)? {
2137 due.push((h.clone(), a.clone(), cp));
2138 }
2139 }
2140 }
2141
2142 let mut out = Vec::new();
2143 for (rec_hash, applied, checkpoint) in due {
2144 let metric = applied.metric.as_ref().unwrap();
2145 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
2146 continue; };
2148 let bound = cost_bound_for(policy, metric);
2149 let base = baseline_at_apply(
2150 sub,
2151 metric,
2152 applied.applied_at_ms,
2153 baseline_kind_for(policy, metric),
2154 bound.map(|b| b.field.as_str()),
2155 )?;
2156 let baseline = base.value;
2157 let tolerance = tolerance_for(policy, metric, &base);
2160 let regressed = crate::recommendation::is_regression(
2161 baseline,
2162 current,
2163 metric.higher_is_better,
2164 tolerance,
2165 );
2166 let (current_run_id, cost) = current_run_and_cost(sub, metric, applied.applied_at_ms, bound, &base)?;
2171 let costlier = cost.as_ref().is_some_and(|c| c.breached());
2172 let verdict = if regressed {
2173 "regressed"
2174 } else if costlier {
2175 "held_costlier"
2176 } else {
2177 "held"
2178 };
2179 p.outcomes.entry(rec_hash.clone()).or_default().push(
2180 crate::recommendation::OutcomeResult {
2181 rec_hash: rec_hash.clone(),
2182 metric: metric.metric.clone(),
2183 baseline,
2184 current,
2185 verdict: verdict.into(),
2186 baseline_kind: base.kind.into(),
2187 baseline_run_id: base.run_id.clone(),
2188 best_before: base.best_before,
2189 tolerance,
2190 current_run_id: current_run_id.clone(),
2191 cost: cost.clone(),
2192 horizon_ms: checkpoint.as_ms().unwrap_or(0),
2195 checkpoint: checkpoint.as_ms().is_none().then_some(checkpoint),
2196 measured_at_ms: now_ms,
2197 },
2198 );
2199 p.measured.entry(rec_hash.clone()).or_default().push(checkpoint);
2200 if regressed || costlier {
2201 out.push(OutcomeInput {
2202 rec_hash,
2203 target_ref: applied.target_ref.clone(),
2204 metric: metric.metric.clone(),
2205 baseline,
2206 current,
2207 unit: metric.unit.clone(),
2208 higher_is_better: metric.higher_is_better,
2209 baseline_kind: base.kind.into(),
2210 baseline_run_id: base.run_id,
2211 best_before: base.best_before,
2212 tolerance,
2213 current_run_id,
2214 cost,
2215 });
2216 }
2217 }
2218 Ok(out)
2219}
2220
2221fn cost_bound_for<'p>(
2223 policy: &'p crate::policy::Policy,
2224 metric: &crate::recommendation::MetricSnapshot,
2225) -> Option<&'p crate::policy::CostBound> {
2226 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2227 (Some(e), Some((hash, _))) if e.hash == hash => e.cost.as_ref(),
2228 _ => None,
2229 }
2230}
2231
2232fn current_run_and_cost<S: SubstrateRead>(
2238 sub: &S,
2239 metric: &crate::recommendation::MetricSnapshot,
2240 applied_at_ms: i64,
2241 bound: Option<&crate::policy::CostBound>,
2242 base: &BaselineRead,
2243) -> Result<(Option<String>, Option<crate::recommendation::CostRead>)> {
2244 let Some((evalset, _)) = crate::eval::parse_evalset_metric(&metric.metric) else {
2245 return Ok((None, None));
2246 };
2247 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(applied_at_ms))? else {
2248 return Ok((None, None));
2249 };
2250 let cost = bound.map(|b| {
2251 let current = crate::eval::run_value(&run, &b.field);
2252 let status = match (base.cost, current) {
2253 (Some(bl), Some(cur)) if cur > bl * b.max_increase_ratio + 1e-9 => "breached",
2254 (Some(_), Some(_)) => "within",
2255 _ => "not_measurable",
2256 };
2257 crate::recommendation::CostRead {
2258 field: b.field.clone(),
2259 max_increase_ratio: b.max_increase_ratio,
2260 baseline: base.cost,
2261 current,
2262 status: status.into(),
2263 }
2264 });
2265 Ok((Some(run.run_id), cost))
2266}
2267
2268fn tolerance_for(
2273 policy: &crate::policy::Policy,
2274 metric: &crate::recommendation::MetricSnapshot,
2275 base: &BaselineRead,
2276) -> f64 {
2277 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2278 (Some(e), Some((hash, field))) if e.hash == hash => e
2279 .min_effect
2280 .map(|m| m.resolve(field, base.total.unwrap_or(metric.n)))
2281 .unwrap_or(0.0),
2282 _ => 0.0,
2283 }
2284}
2285
2286fn checkpoint_due<S: SubstrateRead>(
2295 sub: &S,
2296 metric: &crate::recommendation::MetricSnapshot,
2297 applied_at_ms: i64,
2298 cp: Checkpoint,
2299 now_ms: i64,
2300) -> Result<bool> {
2301 Ok(match cp {
2302 Checkpoint::AfterMs(ms) => now_ms - applied_at_ms >= ms,
2303 Checkpoint::AfterRuns(n) => match crate::eval::parse_evalset_metric(&metric.metric) {
2304 Some((evalset, _)) => {
2305 crate::eval::eval_runs(sub, evalset, Some(applied_at_ms + 1))?.len() >= n as usize
2306 }
2307 None => false,
2308 },
2309 Checkpoint::AfterGrains(n) => count_new(sub, Some(applied_at_ms))?.grains >= n as u64,
2310 })
2311}
2312
2313pub(crate) struct BaselineRead {
2315 pub value: f64,
2316 pub kind: &'static str,
2319 pub run_id: Option<String>,
2320 pub best_before: Option<f64>,
2321 pub total: Option<u64>,
2323 pub cost: Option<f64>,
2326}
2327
2328fn baseline_kind_for(
2331 policy: &crate::policy::Policy,
2332 metric: &crate::recommendation::MetricSnapshot,
2333) -> crate::policy::BaselineKind {
2334 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2335 (Some(e), Some((hash, _))) if e.hash == hash => e.baseline,
2336 _ => crate::policy::BaselineKind::default(),
2337 }
2338}
2339
2340pub(crate) fn baseline_at_apply<S: SubstrateRead>(
2358 sub: &S,
2359 metric: &crate::recommendation::MetricSnapshot,
2360 applied_at_ms: i64,
2361 kind: crate::policy::BaselineKind,
2362 cost_field: Option<&str>,
2363) -> Result<BaselineRead> {
2364 use crate::policy::BaselineKind;
2365 if let Some((evalset, field)) = crate::eval::parse_evalset_metric(&metric.metric) {
2366 let before: Vec<(crate::eval::EvalRun, f64)> = crate::eval::eval_runs(sub, evalset, None)?
2369 .into_iter()
2370 .filter(|r| r.recorded_ms < applied_at_ms)
2371 .filter_map(|r| crate::eval::run_value(&r, field).map(|v| (r, v)))
2372 .collect();
2373 if let Some(newest) = before.last() {
2374 let best = before
2377 .iter()
2378 .fold(None::<&(crate::eval::EvalRun, f64)>, |acc, r| match acc {
2379 None => Some(r),
2380 Some(b) => {
2381 let better = if metric.higher_is_better { r.1 > b.1 } else { r.1 < b.1 };
2382 Some(if better { r } else { b })
2383 }
2384 })
2385 .expect("non-empty");
2386 let pick = match kind {
2387 BaselineKind::NewestBeforeApply => newest,
2388 BaselineKind::HighWater => best,
2389 };
2390 return Ok(BaselineRead {
2391 value: pick.1,
2392 kind: kind.as_str(),
2393 run_id: Some(pick.0.run_id.clone()),
2394 best_before: Some(best.1),
2395 total: Some(pick.0.total()),
2396 cost: cost_field.and_then(|f| crate::eval::run_value(&pick.0, f)),
2397 });
2398 }
2399 }
2400 Ok(BaselineRead { value: metric.baseline, kind: "snapshot", run_id: None, best_before: None, total: None, cost: None })
2401}
2402
2403pub(crate) fn measure_metric<S: SubstrateRead>(
2405 sub: &S,
2406 metric: &crate::recommendation::MetricSnapshot,
2407 since_ms: i64,
2408) -> Result<Option<f64>> {
2409 match metric.metric.as_str() {
2410 "tool_error_recurrence" => {
2415 let Some(tool) = &metric.subject else { return Ok(None) };
2416 let tools = sub.grains_of_type(
2417 crate::model::grain_type::TOOL,
2418 None,
2419 ReadOpts { live_only: true, since_ms: Some(since_ms) },
2420 )?;
2421 let n = tools
2422 .iter()
2423 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
2424 .filter(|t| {
2425 metric.relation.as_deref().is_none_or(|sig| {
2428 crate::analyzers::tool_failure::normalize_signature(
2429 t.tool_content().unwrap_or(""),
2430 ) == sig
2431 })
2432 })
2433 .count();
2434 Ok(Some(n as f64))
2435 }
2436 "contradiction_recurrence" => {
2440 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
2441 return Ok(None);
2442 };
2443 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
2444 let distinct: BTreeSet<String> = facts
2445 .iter()
2446 .filter(|f| {
2447 f.fact_relation()
2448 .is_some_and(|r| normalize_ident(r) == *relation)
2449 })
2450 .filter_map(|f| f.fact_object().map(normalize_ident))
2451 .collect();
2452 Ok(Some(distinct.len().saturating_sub(1) as f64))
2453 }
2454 m if m.starts_with("evalset:") => {
2467 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
2468 return Ok(None);
2469 };
2470 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
2471 return Ok(None);
2472 };
2473 Ok(crate::eval::run_value(&run, field))
2474 }
2475 _ => Ok(None),
2476 }
2477}
2478
2479fn scoped_live_facts<S: SubstrateRead>(
2482 sub: &S,
2483 namespace: Option<&str>,
2484 subject: &str,
2485) -> Result<Vec<GrainRecord>> {
2486 let facts = sub.grains_of_type(
2487 crate::model::grain_type::FACT,
2488 None,
2489 ReadOpts { live_only: true, since_ms: None },
2490 )?;
2491 Ok(facts
2492 .into_iter()
2493 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
2494 .filter(|f| {
2495 f.fact_subject()
2496 .is_some_and(|s| normalize_ident(s) == subject)
2497 })
2498 .collect())
2499}
2500
2501fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
2513 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
2514 return false;
2515 };
2516 if fields.is_empty() {
2517 return false;
2518 }
2519 let Ok(Some(grain)) = sub.grain(&target) else {
2520 return false;
2521 };
2522 if grain.valid_to_ms.is_some() {
2532 return false;
2533 }
2534 fields.iter().all(|(k, v)| {
2536 if k == "namespace" {
2537 return v
2538 .as_str()
2539 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2540 }
2541 match (v, grain.fields.get(k)) {
2542 (Value::String(a), Some(Value::String(b))) => {
2543 normalize_ident(a) == normalize_ident(b)
2544 }
2545 (a, Some(b)) => a == b,
2546 (_, None) => false,
2547 }
2548 })
2549}
2550
2551const EVIDENCE_CAP: usize = 64;
2561const CITED_SEED_CAP: usize = 24;
2562const TOOL_SEED_CAP: usize = 16;
2563const NOTE_SEED_CAP: usize = 8;
2568const HARNESS_EVIDENCE_KINDS: &[&str] = &["fold_summary"];
2587const HARNESS_SEED_CAP: usize = 6;
2588const LENS_RESERVE: usize = 24;
2589
2590const MIN_LLM_CONFIDENCE: f64 = 0.75;
2593
2594macro_rules! discover_instructions {
2604 ($scoring:literal) => {
2605 concat!(
2606 "You review an agent's memory for quality. \
2607Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2608findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2609stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2610recurring cost or hand-off the agent's own setup could remove). \
2611The deterministic findings already cover what the ERROR TEXT says; restating \
2612one of them earns nothing. The evidence may also contain OUTCOME records — a \
2613run's observable shape together with whether it was accepted or rejected. A \
2614problem that raised no error at all is exactly the kind the deterministic \
2615checks cannot see, so compare the rejected outcomes against the accepted \
2616ones: a feature they share and the accepted ones lack is a candidate rule. \
2617Require at least two rejected outcomes before proposing one — a single \
2618rejection is an anecdote, not a pattern. ",
2619 $scoring,
2620 " The 'approved' and 'rejected' lists, when \
2621present, show findings this reviewer recently accepted or rejected — prefer the \
2622kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2623or more evidence items from the bundle by their 'id' (or 'hash'), name a \
2624'target', and include your confidence 0.0-1.0. Return JSON: \
2625{\"recommendations\":[{\"summary\":\"...\",\
2626\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<id>\"],\
2627\"confidence\":0.0,\"proposal\":{...}}]}. \
2628OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2629you are not asking to change anything. Include it ONLY when the evidence \
2630supports a specific change, choosing exactly one kind: \
2631(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2632\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
2633an action the agent itself takes on the next occasion. Either ADD an action \
2634it is failing to take ('Record the vendor name and the amount on every \
2635invoice, not just the date') or ORDER one that goes wrong ('Refund a \
2636subscription before cancelling it; refunds on cancelled subscriptions are \
2637refused'). Name the action, not a check on it: 'validate', 'verify' and \
2638'ensure ... is correct' describe a review step the agent has no way to \
2639perform, and such a rule changes nothing even once applied. \
2640(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2641entity target — a durable fact the agent keeps having to be told (an alias, a \
2642settled default, a preference). 'relation' is a short identifier (letters, \
2643digits, _ - . :), not a sentence. \
2644(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2645\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2646assembles the agent's context, when the evidence shows it retrieves the wrong \
2647things. Give the FULL new body; it replaces the old one. \
2648(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2649with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2650workflow. Only these paths are editable: 'edges.<i>.cond', \
2651'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2652holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2653(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2654— full replacement source for that tool. It is applied only after a recorded \
2655evaluation run passes, so propose one only when the evidence shows the current \
2656code is the defect. \
2657The subject of a fact, the name of a query, the plan hash and the tool name \
2658all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2659becomes a change a human reviewer may apply, so it must be fully supported by \
2660the cited evidence. Propose nothing you cannot ground in the evidence."
2661 )
2662 };
2663}
2664
2665const DISCOVER_INSTRUCTIONS: &str = discover_instructions!(
2667 "SCORING: propose a finding ONLY if you \
2668are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2669useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2670earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2671when there is nothing worth flagging."
2672);
2673
2674const DISCOVER_LEARNER_INSTRUCTIONS: &str = discover_instructions!(
2676 "SCORING: you are the learning stage of a deployed agent, and what you \
2677propose now is what it will do differently next time — a lesson you withhold \
2678is a mistake it repeats. A correct, actionable proposal earns 1; a wrong or \
2679trivial one is penalized 1; returning nothing while the evidence holds a \
2680recurring failure, two or more rejected outcomes, an instruction from a \
2681person, or a multi-step procedure the agent completed successfully that no \
2682saved skill or plan covers, is ALSO penalized 1. Abstain only when the \
2683evidence shows none of those. Prefer the one proposal that addresses the most \
2684frequent or most costly failure — or, when nothing failed, the procedure that \
2685worked — over several speculative ones, and report your confidence honestly — \
2686an independent verifier, not you, decides what survives."
2687);
2688
2689const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2695fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2696the facts it relies on are actually present in the cited evidence, NOT that its \
2697conclusion is stated verbatim. Decompose the finding into the factual claims it \
2698depends on. Mark supported=true when those facts are present in the evidence \
2699(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2700on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2701different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2702
2703const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2705each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2706never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2707SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2708'possible' findings with no concrete defect, and reject any claimed \
2709inconsistency or contradiction that is not backed by at least two actually \
2710conflicting facts in the cited evidence. (2) Context — does the finding \
2711correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2712finding when it names a genuine, specific problem grounded in its evidence and \
2713materially useful to a human reviewer; otherwise reject it, and default to \
2714keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2715redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2716grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2717KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2718{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2719
2720const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2722guidance note to help a human reviewer decide. Do not restate the finding. Return \
2723JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2724
2725fn push_evidence(
2731 evidence: &mut Vec<crate::llm::EvidenceItem>,
2732 bundle: &mut BTreeSet<String>,
2733 ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2734 g: &GrainRecord,
2735 attribution: crate::policy::EvidenceAttribution,
2736) {
2737 if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2738 ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2739 evidence.push(crate::llm::EvidenceItem {
2740 id: format!("e{}", evidence.len() + 1),
2741 hash: g.hash.clone(),
2742 grain_type: g.grain_type.clone(),
2743 text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
2744 });
2745 }
2746}
2747
2748pub(crate) fn resolve_citation(
2756 cite: &str,
2757 bundle: &BTreeSet<String>,
2758 id_to_hash: &std::collections::BTreeMap<&str, &str>,
2759) -> Option<String> {
2760 let cite = cite.trim();
2761 if bundle.contains(cite) {
2762 return Some(cite.to_string());
2763 }
2764 if let Some(h) = id_to_hash.get(cite) {
2765 return Some((*h).to_string());
2766 }
2767 const MIN_PREFIX: usize = 12;
2768 if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
2769 let lower = cite.to_ascii_lowercase();
2770 let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
2771 if let (Some(h), None) = (it.next(), it.next()) {
2772 return Some(h.clone());
2773 }
2774 }
2775 None
2776}
2777
2778fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
2781 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2782 return format!("{s} {r} {o}");
2783 }
2784 if let Some(t) = g.tool_name() {
2789 let status = if g.is_error() { "error" } else { "ok" };
2790 let out = g.tool_content().unwrap_or("");
2791 let input = match g.fields.get("input") {
2797 Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
2798 Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
2799 _ => String::new(),
2800 };
2801 return format!("tool {t}{input} {status}: {out}");
2802 }
2803 for key in ["content", "body", "text", "summary", "object"] {
2811 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2812 if v.is_empty() {
2813 continue;
2814 }
2815 if attribution == crate::policy::EvidenceAttribution::Anonymous {
2826 return v.to_string();
2827 }
2828 if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
2829 if !who.is_empty() {
2830 let kind = g
2831 .fields
2832 .get("observer_type")
2833 .and_then(|v| v.as_str())
2834 .unwrap_or("");
2835 let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
2836 let mut prefix = if kind == "human" {
2837 format!("{who} (a person) said")
2838 } else {
2839 format!("{who} observed")
2840 };
2841 if !about.is_empty() {
2842 prefix.push_str(&format!(" of {about}"));
2843 }
2844 return format!("{prefix}: {v}");
2845 }
2846 }
2847 return v.to_string();
2848 }
2849 }
2850 String::new()
2851}
2852
2853fn sanitize_lesson(s: &str) -> String {
2858 sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2859}
2860
2861fn sanitize_line(s: &str, max: usize) -> String {
2867 let cleaned: String =
2868 s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2869 crate::llm::cap(cleaned.trim(), max)
2870}
2871
2872fn sanitize_relation(s: &str) -> Option<String> {
2876 let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2877 if r.is_empty()
2878 || !r
2879 .chars()
2880 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2881 {
2882 return None;
2883 }
2884 Some(r)
2885}
2886
2887fn safe_definition_body(body: &str) -> bool {
2903 if body.contains('{') || body.contains('}') {
2904 return false;
2905 }
2906 !body
2907 .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2908 .any(|tok| {
2909 ["FORGET", "PURGE", "DROP", "DEFINE"]
2910 .iter()
2911 .any(|kw| tok.eq_ignore_ascii_case(kw))
2912 })
2913}
2914
2915fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2920 let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2921 match resolved {
2922 Some(r) => format!("{summary} {}", r.rendered),
2923 None => summary,
2924 }
2925}
2926
2927struct ValidatedDraft {
2933 draft: crate::llm::LlmDraft,
2934 target_ref: String,
2935 cited: Vec<String>,
2936 resolved: Option<ResolvedProposal>,
2937}
2938
2939struct ResolvedProposal {
2943 action: ActionKind,
2944 proposal: Proposal,
2945 rendered: String,
2948 summary_key: &'static str,
2949 summary_args: serde_json::Map<String, Value>,
2950 rollbackable: bool,
2951 evalset_hash: Option<String>,
2952 importance: f64,
2953 fact_fields: Option<serde_json::Map<String, Value>>,
2957 extra_statements: Vec<String>,
2960 replay: Option<Value>,
2963}
2964
2965fn plan_edit_allowed(path: &str) -> bool {
2975 let seg: Vec<&str> = path.split('.').collect();
2976 match seg.as_slice() {
2977 ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2978 ["retries", node] => !node.is_empty(),
2979 _ => false,
2980 }
2981}
2982
2983fn plan_get(body: &Value, path: &str) -> Value {
2986 let mut cur = body;
2987 for seg in path.split('.') {
2988 cur = match cur {
2989 Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2990 Some(v) => v,
2991 None => return Value::Null,
2992 },
2993 Value::Object(o) => match o.get(seg) {
2994 Some(v) => v,
2995 None => return Value::Null,
2996 },
2997 _ => return Value::Null,
2998 };
2999 }
3000 cur.clone()
3001}
3002
3003fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
3009 let segs: Vec<&str> = path.split('.').collect();
3010 let Some((last, parents)) = segs.split_last() else {
3011 return false;
3012 };
3013 let mut cur = body;
3014 for (depth, seg) in parents.iter().enumerate() {
3015 cur = match cur {
3016 Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
3017 Some(v) => v,
3018 None => return false,
3019 },
3020 Value::Object(o) => {
3021 if depth == 0 && *seg == "retries" && !o.contains_key("retries") {
3022 o.insert("retries".into(), Value::Object(serde_json::Map::new()));
3023 }
3024 match o.get_mut(*seg) {
3025 Some(v) => v,
3026 None => return false,
3027 }
3028 }
3029 _ => return false,
3030 };
3031 }
3032 match cur {
3033 Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
3034 Some(slot) => {
3035 *slot = to;
3036 true
3037 }
3038 None => false,
3039 },
3040 Value::Object(o) => {
3041 o.insert((*last).to_string(), to);
3042 true
3043 }
3044 _ => false,
3045 }
3046}
3047
3048fn plan_value_ok(path: &str, to: &Value) -> bool {
3053 let seg: Vec<&str> = path.split('.').collect();
3054 match seg.as_slice() {
3055 ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
3056 !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
3057 }),
3058 ["edges", _, "max_cycles"] | ["retries", _] => {
3059 to.as_u64().is_some_and(|n| n <= 1_000)
3060 }
3061 _ => false,
3062 }
3063}
3064
3065fn resolve_proposal<S: OmsSubstrate>(
3071 sub: &S,
3072 d: &crate::llm::LlmDraft,
3073 target: &TargetRef,
3074 cited: &[String],
3075 ns_by_hash: &std::collections::BTreeMap<String, String>,
3076 caps: Capabilities,
3077 policy: &crate::policy::Policy,
3078) -> Option<ResolvedProposal> {
3079 use crate::llm::DraftProposal as P;
3080 let (skills, plans) = (&policy.skills, &policy.plans);
3081 let mut args = serde_json::Map::new();
3082 match d.parsed_proposal()? {
3083 P::Plan { description, when_to_use, nodes, edges } => {
3085 if !plans.enabled || !caps.plans {
3086 return None;
3087 }
3088 let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
3089 derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
3090 args.insert("name".into(), Value::from(name.clone()));
3091 args.insert("nodes".into(), Value::from(n_nodes as u64));
3092 let mut stmts = vec![match &existing_skill {
3093 Some(h) => cal::supersede(h, "skill", &skill),
3094 None => cal::add("skill", &skill),
3095 }];
3096 let (summary_key, kind) = match &workflow {
3103 Some(wf) => {
3104 stmts.push(match &existing_plan {
3105 Some(h) => cal::supersede(h, "workflow", wf),
3106 None => cal::add("workflow", wf),
3107 });
3108 args.insert("edges".into(), Value::from(n_edges as u64));
3109 ("llm.plan", "plan")
3110 }
3111 None => {
3112 args.insert("steps".into(), Value::from(n_nodes as u64));
3113 ("llm.skill", "skill")
3114 }
3115 };
3116 let patched = existing_skill.is_some() || (workflow.is_some() && existing_plan.is_some());
3117 let (action, verb) = if patched {
3118 (ActionKind::Revise, "revise")
3119 } else {
3120 (ActionKind::Record, "record")
3121 };
3122 Some(ResolvedProposal {
3123 action,
3124 proposal: Proposal::Cal { cal: cal::batch(&stmts) },
3125 rendered: format!(
3126 "Proposed {kind} to {verb}: \"{name}\" — {n_nodes} steps; when: {}",
3127 skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
3128 ),
3129 summary_key,
3130 summary_args: args,
3131 rollbackable: true,
3132 evalset_hash: None,
3133 importance: 0.65,
3134 fact_fields: None,
3135 extra_statements: Vec::new(),
3136 replay: None,
3137 })
3138 }
3139 P::Skill { description, when_to_use, steps } => {
3141 if !skills.enabled {
3142 return None;
3143 }
3144 let SkillFields { fields, name, n_steps, existing } =
3145 derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
3146 args.insert("name".into(), Value::from(name.clone()));
3147 args.insert("steps".into(), Value::from(n_steps as u64));
3148 let (action, cal, verb) = match existing {
3151 Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
3152 None => (ActionKind::Record, cal::add("skill", &fields), "record"),
3153 };
3154 Some(ResolvedProposal {
3155 action,
3156 proposal: Proposal::Cal { cal },
3157 rendered: format!(
3158 "Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
3159 fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
3160 ),
3161 summary_key: "llm.skill",
3162 summary_args: args,
3163 rollbackable: true,
3164 evalset_hash: None,
3165 importance: 0.6,
3166 fact_fields: None,
3167 extra_statements: Vec::new(),
3168 replay: None,
3169 })
3170 }
3171 P::Consolidation { lesson, supersedes } => {
3173 let lesson = sanitize_lesson(&lesson);
3174 if lesson.is_empty() {
3175 return None;
3176 }
3177 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
3178 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("").to_string();
3179 let mut hashes: Vec<String> = supersedes.into_iter().collect();
3184 hashes.sort();
3185 hashes.dedup();
3186 let mut members = Vec::new();
3187 for h in &hashes {
3188 let g = sub.grain(h).ok().flatten()?;
3189 if !g.is_live()
3190 || g.fact_relation() != Some("lesson")
3191 || g.fact_subject().is_none_or(|s| normalize_ident(s) != normalize_ident(&subject))
3192 {
3193 return None;
3194 }
3195 members.push(g);
3196 }
3197 if members.len() < 2 || members.iter().any(|m| m.namespace != members[0].namespace) {
3198 return None;
3199 }
3200 let ns = members[0].namespace.clone();
3201 let mut fields = fields;
3202 if !ns.is_empty() {
3203 fields.insert("namespace".into(), Value::from(ns.clone()));
3204 }
3205 fields.insert("consolidates".into(), Value::from(hashes.clone()));
3206 let extra_statements: Vec<String> = members
3211 .iter()
3212 .map(|m| {
3213 let mut marker = serde_json::Map::new();
3214 marker.insert("subject".into(), Value::from(subject.clone()));
3215 marker.insert("relation".into(), Value::from("mg:lesson_consolidated"));
3216 marker.insert("object".into(), Value::from(lesson.clone()));
3217 if !ns.is_empty() {
3218 marker.insert("namespace".into(), Value::from(ns.clone()));
3219 }
3220 cal::supersede(&m.hash, "fact", &marker)
3221 })
3222 .collect();
3223 args.insert("lesson".into(), Value::from(lesson.clone()));
3224 args.insert("count".into(), Value::from(members.len() as u64));
3225 Some(ResolvedProposal {
3226 action: ActionKind::Consolidate,
3227 proposal: Proposal::Cal { cal: cal::batch(&extra_statements) },
3228 rendered: format!(
3229 "Proposed consolidation of {} lessons on \"{subject}\" into one: \"{lesson}\"",
3230 members.len()
3231 ),
3232 summary_key: "llm.consolidation",
3233 summary_args: args,
3234 rollbackable: true,
3235 evalset_hash: None,
3236 importance: 0.6,
3237 fact_fields: Some(fields),
3238 extra_statements,
3239 replay: None,
3240 })
3241 }
3242 P::Lesson { lesson } => {
3244 let lesson = sanitize_lesson(&lesson);
3245 if lesson.is_empty() {
3246 return None;
3247 }
3248 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
3249 args.insert("lesson".into(), Value::from(lesson.clone()));
3250 Some(ResolvedProposal {
3251 action: ActionKind::ClusterFailure,
3255 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
3256 rendered: format!("Proposed lesson to record: \"{lesson}\""),
3257 summary_key: "llm.lesson",
3258 summary_args: args,
3259 rollbackable: true,
3260 evalset_hash: None,
3261 importance: 0.5,
3262 fact_fields: Some(fields),
3263 extra_statements: Vec::new(),
3264 replay: None,
3265 })
3266 }
3267 P::Fact { relation, object } => {
3269 let relation = sanitize_relation(&relation)?;
3270 let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
3271 if object.is_empty() {
3272 return None;
3273 }
3274 let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
3275 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
3276 args.insert("relation".into(), Value::from(relation.clone()));
3277 args.insert("object".into(), Value::from(object.clone()));
3278 Some(ResolvedProposal {
3279 action: ActionKind::Record,
3280 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
3281 rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
3282 summary_key: "llm.fact",
3283 summary_args: args,
3284 rollbackable: true,
3285 evalset_hash: None,
3286 importance: 0.5,
3287 fact_fields: Some(fields),
3288 extra_statements: Vec::new(),
3289 replay: None,
3290 })
3291 }
3292 P::QueryRevision { body } => {
3294 let name = target.opaque();
3295 if name.is_empty()
3299 || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
3300 {
3301 return None;
3302 }
3303 let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
3304 if body.is_empty() || !safe_definition_body(&body) {
3305 return None;
3306 }
3307 let stmt = match target.scheme() {
3308 "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
3309 "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
3310 _ => return None,
3311 };
3312 sub.validate_cal(&stmt).ok()?;
3317 sub.definition_inverse(&stmt).ok().flatten()?;
3318 args.insert("name".into(), Value::from(name));
3319 args.insert("body".into(), Value::from(body.clone()));
3320 Some(ResolvedProposal {
3321 action: ActionKind::Revise,
3322 proposal: Proposal::Cal { cal: stmt },
3323 rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
3324 summary_key: "llm.query_revision",
3325 summary_args: args,
3326 rollbackable: true,
3327 evalset_hash: None,
3328 importance: 0.6,
3329 fact_fields: None,
3330 extra_statements: Vec::new(),
3331 replay: None,
3332 })
3333 }
3334 P::PlanRevision { edits } => {
3336 if !caps.plans
3337 || target.scheme() != "grain"
3338 || edits.is_empty()
3339 || edits.len() > crate::llm::MAX_PLAN_EDITS
3340 {
3341 return None;
3342 }
3343 let hash = target.opaque();
3344 let g = sub.grain(hash).ok().flatten()?;
3345 if g.grain_type != "workflow" || !g.is_live() {
3346 return None;
3347 }
3348 let mut body = Value::Object(g.fields.clone());
3349 let mut deltas = Vec::new();
3350 let nodes: std::collections::BTreeSet<String> = body
3351 .get("nodes")
3352 .and_then(Value::as_array)
3353 .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
3354 .unwrap_or_default();
3355 for e in &edits {
3356 if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
3357 return None;
3358 }
3359 if let Some(node) = e.path.strip_prefix("retries.") {
3364 if !nodes.contains(node) {
3365 return None;
3366 }
3367 }
3368 if plan_get(&body, &e.path) != e.from {
3371 return None;
3372 }
3373 if e.from == e.to {
3377 return None;
3378 }
3379 if !plan_set(&mut body, &e.path, e.to.clone()) {
3380 return None;
3381 }
3382 deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
3383 }
3384 sub.validate_plan(&body).ok()?;
3388 let replay = sub.plan_replay(hash, &body).ok().flatten();
3397 let refused = match (&policy.plan_replay, &replay) {
3398 (Some(gate), Some(report)) => gate.refusal(report),
3399 _ => None,
3400 };
3401 let Value::Object(fields) = body else {
3402 return None;
3403 };
3404 args.insert("plan".into(), Value::from(hash));
3405 args.insert("edits".into(), Value::from(deltas.join("; ")));
3406 if let Some(reason) = refused {
3407 let mut data = serde_json::Map::new();
3408 data.insert("plan".into(), Value::from(hash));
3409 data.insert("edits".into(), Value::from(deltas.clone()));
3410 data.insert("refused".into(), Value::from(reason.clone()));
3411 args.insert("reason".into(), Value::from(reason.clone()));
3412 return Some(ResolvedProposal {
3413 action: ActionKind::Flag,
3414 proposal: Proposal::Data { data },
3415 rendered: format!(
3416 "Plan revision ({}) refused by the rehearsal: {reason}",
3417 deltas.join("; ")
3418 ),
3419 summary_key: "llm.plan_revision_refused",
3420 summary_args: args,
3421 rollbackable: false,
3422 evalset_hash: None,
3423 importance: 0.4,
3424 fact_fields: None,
3425 extra_statements: Vec::new(),
3426 replay,
3427 });
3428 }
3429 let stmt = cal::supersede(hash, "workflow", &fields);
3430 sub.validate_cal(&stmt).ok()?;
3435 Some(ResolvedProposal {
3436 action: ActionKind::Revise,
3437 proposal: Proposal::Cal { cal: stmt },
3438 rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
3439 summary_key: "llm.plan_revision",
3440 summary_args: args,
3441 rollbackable: true,
3442 evalset_hash: None,
3443 importance: 0.7,
3444 fact_fields: None,
3445 extra_statements: Vec::new(),
3446 replay,
3447 })
3448 }
3449 P::CodeRevision { source } => {
3451 if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
3452 return None;
3453 }
3454 if source.chars().count() > crate::llm::MAX_CODE_LEN {
3455 return None;
3456 }
3457 let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
3461 let mut data = serde_json::Map::new();
3462 data.insert("tool".into(), Value::from(target.opaque()));
3463 data.insert("source".into(), Value::from(source.clone()));
3464 args.insert("tool".into(), Value::from(target.opaque()));
3465 args.insert("bytes".into(), Value::from(source.len() as u64));
3466 Some(ResolvedProposal {
3467 action: ActionKind::CodeRevision,
3468 proposal: Proposal::Data { data },
3469 rendered: format!(
3470 "Proposed new source for tool {} ({} bytes), gated by evalset {}",
3471 target.opaque(),
3472 source.len(),
3473 evalset
3474 ),
3475 summary_key: "llm.code_revision",
3476 summary_args: args,
3477 rollbackable: true,
3478 evalset_hash: Some(evalset),
3479 importance: 0.8,
3480 fact_fields: None,
3481 extra_statements: Vec::new(),
3482 replay: None,
3483 })
3484 }
3485 }
3486}
3487
3488#[allow(clippy::too_many_arguments)]
3498fn stamp_llm(
3499 model: &str,
3500 d: &crate::llm::LlmDraft,
3501 target_ref: String,
3502 cited: Vec<String>,
3503 resolved: Option<ResolvedProposal>,
3504 confidence: f64,
3505 now_ms: i64,
3506 scope: &[String],
3507) -> Recommendation {
3508 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3509 let guidance = if d.guidance.trim().is_empty() {
3510 None
3511 } else {
3512 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
3513 };
3514 let replay = resolved.as_ref().and_then(|r| r.replay.clone());
3515 let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
3516 Some(mut r) => {
3517 let content = match &r.fact_fields {
3521 Some(fields) => format!(
3522 "{} {}",
3523 fields.get("relation").and_then(Value::as_str).unwrap_or(""),
3524 fields.get("object").and_then(Value::as_str).unwrap_or("")
3525 ),
3526 None => match &r.proposal {
3527 Proposal::Cal { cal } => cal.clone(),
3528 Proposal::Data { data } => Value::Object(data.clone()).to_string(),
3529 Proposal::Edit { diff, .. } => diff.clone(),
3530 },
3531 };
3532 if let Some(mut fields) = r.fact_fields.take() {
3535 fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
3536 let mut statements = vec![cal::add("fact", &fields)];
3537 statements.extend(r.extra_statements.iter().cloned());
3538 r.proposal = Proposal::Cal { cal: cal::batch(&statements) };
3539 }
3540 let mut args = r.summary_args;
3541 args.insert("text".into(), Value::from(summary_text));
3542 (
3543 r.action,
3544 r.proposal,
3545 Summary::new(r.summary_key, args),
3546 r.rollbackable,
3547 r.importance,
3548 r.evalset_hash,
3549 Some(content),
3550 )
3551 }
3552 None => {
3553 let mut args = serde_json::Map::new();
3554 args.insert("text".into(), Value::from(summary_text));
3555 let mut data = serde_json::Map::new();
3556 data.insert("source".into(), Value::from("llm"));
3557 (
3558 ActionKind::Flag,
3559 Proposal::Data { data },
3560 Summary::new("llm.discover", args),
3561 false,
3562 0.3,
3563 None,
3564 None,
3565 )
3566 }
3567 };
3568 let dedup = match &content {
3572 Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
3573 None => dedup_key("llm", &target_ref, action),
3574 };
3575 Recommendation {
3576 hash: String::new(),
3577 analyzer: "loop.llm/1".to_string(),
3578 params_snapshot: serde_json::Map::new(),
3579 origin: Origin::Llm { model: model.to_string() },
3580 target_ref: target_ref.clone(),
3581 action_kind: action,
3582 dedup_key: dedup,
3583 summary,
3584 severity: Severity::Low,
3585 proposal,
3586 destructive: false,
3587 rollbackable,
3588 evidence: cited,
3589 evidence_query: None,
3590 metric: None,
3591 confidence: confidence.clamp(0.0, 1.0),
3593 importance,
3594 created_at_ms: now_ms,
3595 guidance,
3596 evalset_hash,
3597 near_duplicate_of: Vec::new(),
3598 replay,
3599 scope: crate::recommendation::normalize_scope(scope),
3603 status: RecStatus::Pending,
3604 }
3605}
3606
3607fn skill_instructions(min_steps: u32) -> String {
3617 format!(
3618 " (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
3619\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
3620PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
3621calls that reached its goal, which a later session facing the same situation \
3622should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
3623the tool called and the values that mattered (the field checked, the tag set, the \
3624exact format produced), a one-line description, and 'when_to_use' — the situation \
3625that should trigger it. The skill-name is a short identifier (letters, digits, \
3626_ -). If a saved skill already covers this procedure, use ITS name so it is \
3627patched rather than duplicated. Do not propose a skill for a procedure that \
3628failed, or for one already saved and unchanged. A finding that itself describes \
3629two or more steps the agent should carry out in order ('after listing the \
3630tickets, fetch each, then …') IS a procedure: propose it as a skill or a plan, \
3631never as a lesson — a lesson is one rule, and a procedure written as one is a \
3632procedure nobody can open.",
3633 crate::llm::MAX_SKILL_STEPS
3634 )
3635}
3636
3637fn plan_instructions(min_nodes: u32) -> String {
3641 format!(
3642 " (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
3643\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
3644\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
3645with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
3646but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
3647'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
3648the cited evidence — and what the step does with it; and 'edges' from step to \
3649step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
3650'path != literal', 'path exists' or '!path', where path is dotted names and the \
3651literal is a JSON string, number, true, false or null — no other operators; state \
3652a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
3653earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
3654has branches or a loop; prefer a skill when it is a straight list. If a saved \
3655plan already covers this procedure, use ITS name so it is patched.",
3656 crate::llm::MAX_PLAN_NODES
3657 )
3658}
3659
3660struct PlanFields {
3664 skill: serde_json::Map<String, Value>,
3665 workflow: Option<serde_json::Map<String, Value>>,
3669 name: String,
3670 n_nodes: usize,
3671 n_edges: usize,
3672 existing_skill: Option<String>,
3673 existing_plan: Option<String>,
3674}
3675
3676#[allow(clippy::too_many_arguments)]
3686fn derived_plan_fields<S: SubstrateRead>(
3687 sub: &S,
3688 target: &TargetRef,
3689 description: &str,
3690 when_to_use: &str,
3691 nodes: &[crate::llm::PlanNodeDraft],
3692 edges: &[crate::llm::PlanEdgeDraft],
3693 cited: &[String],
3694 ns_by_hash: &std::collections::BTreeMap<String, String>,
3695 plans: &crate::policy::PlanAuthoring,
3696) -> Option<PlanFields> {
3697 if target.scheme() != "entity" {
3698 return None;
3699 }
3700 let name = sanitize_skill_name(
3701 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3702 )?;
3703 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3704 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3705 if description.is_empty() || when_to_use.is_empty() {
3706 return None;
3707 }
3708 if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
3709 return None;
3710 }
3711 let known_tools: BTreeSet<String> = cited
3713 .iter()
3714 .filter_map(|h| sub.grain(h).ok().flatten())
3715 .filter_map(|g| g.tool_name().map(normalize_ident))
3716 .collect();
3717 let mut ids: Vec<String> = Vec::new();
3718 let mut steps: Vec<String> = Vec::new();
3719 let mut seen: BTreeSet<String> = BTreeSet::new();
3720 let mut grounded = 0usize;
3721 for n in nodes {
3722 let id = sanitize_skill_name(&n.id)?;
3723 if !seen.insert(id.clone()) {
3724 return None; }
3726 let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
3727 if step.is_empty() {
3728 return None;
3729 }
3730 let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
3740 if !tool.is_empty() && known_tools.contains(&normalize_ident(&tool)) {
3741 grounded += 1;
3742 steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
3743 } else {
3744 steps.push(format!("{}. {id}: {step}", steps.len() + 1));
3745 }
3746 ids.push(id);
3747 }
3748 if grounded == 0 {
3751 return None;
3752 }
3753 let mut edge_vals: Vec<Value> = Vec::new();
3758 let mut flow_lines: Vec<String> = Vec::new();
3759 let mut runnable = true;
3760 for e in edges {
3761 let (Some(src), Some(dst)) = (sanitize_skill_name(&e.src), sanitize_skill_name(&e.dst)) else {
3762 runnable = false;
3763 continue;
3764 };
3765 if !seen.contains(&src) || !seen.contains(&dst) {
3766 flow_lines.push(format!("{} → {}", e.src.trim(), e.dst.trim()));
3768 runnable = false;
3769 continue;
3770 }
3771 let mut ev = serde_json::Map::new();
3772 ev.insert("src".into(), Value::from(src.clone()));
3773 ev.insert("dst".into(), Value::from(dst.clone()));
3774 let mut label = format!("{src} → {dst}");
3775 if let Some(c) = e
3776 .cond
3777 .as_deref()
3778 .map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
3779 .filter(|c| !c.is_empty())
3780 {
3781 label.push_str(&format!(" if {c}"));
3782 ev.insert("cond".into(), Value::from(c));
3783 }
3784 if let Some(m) = e.max_cycles {
3785 if m == 0 || m > 100 {
3786 runnable = false;
3787 } else {
3788 label.push_str(&format!(" (at most {m} times)"));
3789 ev.insert("max_cycles".into(), Value::from(m));
3790 }
3791 }
3792 flow_lines.push(label);
3793 edge_vals.push(Value::Object(ev));
3794 }
3795 if edge_vals.len() > 4 * ids.len() {
3796 runnable = false;
3797 }
3798 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3800 for h in cited {
3801 if let Some(ns) = ns_by_hash.get(h) {
3802 if !ns.is_empty() {
3803 *ns_counts.entry(ns.as_str()).or_default() += 1;
3804 }
3805 }
3806 }
3807 let ns = ns_counts
3808 .iter()
3809 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3810 .map(|(ns, _)| ns.to_string());
3811
3812 let mut workflow = serde_json::Map::new();
3815 workflow.insert("nodes".into(), Value::from(ids.clone()));
3816 workflow.insert("edges".into(), Value::Array(edge_vals));
3817 workflow.insert("name".into(), Value::from(name.clone()));
3818 if let Some(ns) = &ns {
3819 workflow.insert("namespace".into(), Value::from(ns.clone()));
3820 }
3821 let workflow = (runnable && sub.validate_plan(&Value::Object(workflow.clone())).is_ok())
3825 .then_some(workflow);
3826
3827 let mut instructions = steps.join("\n");
3830 if !flow_lines.is_empty() {
3831 instructions.push_str("\n\nFlow:\n");
3832 instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
3833 }
3834 let mut skill = serde_json::Map::new();
3835 skill.insert("name".into(), Value::from(name.clone()));
3836 skill.insert("description".into(), Value::from(description));
3837 skill.insert("when_to_use".into(), Value::from(when_to_use));
3838 skill.insert("instructions".into(), Value::from(instructions));
3839 if let Some(ns) = &ns {
3840 skill.insert("namespace".into(), Value::from(ns.clone()));
3841 }
3842 let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
3843 sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
3844 .ok()?
3845 .into_iter()
3846 .find(|g| pick(g))
3847 .map(|g| g.hash)
3848 };
3849 let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
3850 let existing_plan = workflow
3851 .is_some()
3852 .then(|| live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str())))
3853 .flatten();
3854 let n_edges = workflow
3855 .as_ref()
3856 .and_then(|w| w.get("edges"))
3857 .and_then(Value::as_array)
3858 .map_or(0, |a| a.len());
3859 Some(PlanFields { skill, workflow, name, n_nodes: ids.len(), n_edges, existing_skill, existing_plan })
3860}
3861
3862pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
3864
3865fn detect_premise_drift<S: OmsSubstrate>(
3881 sub: &S,
3882 p: &mut LoopPersisted,
3883 now_ms: i64,
3884) -> Result<Vec<OutcomeInput>> {
3885 let mut out = Vec::new();
3886 let applied: Vec<(String, String, Vec<String>)> = p
3887 .applied
3888 .iter()
3889 .filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
3890 .map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
3891 .collect();
3892 for (rec_hash, target_ref, own) in applied {
3893 let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
3894 if rec.evidence.is_empty() {
3895 continue;
3896 }
3897 let mut moved = 0u64;
3898 for e in &rec.evidence {
3899 match sub.grain(e)? {
3900 None => moved += 1, Some(g) => {
3902 let Some(newer) = &g.superseded_by else { continue };
3903 if own.iter().any(|c| c == newer) {
3904 continue; }
3906 match sub.grain(newer)? {
3907 None => moved += 1,
3910 Some(n) => {
3911 if !same_value(&g, &n) {
3912 moved += 1;
3913 }
3914 }
3915 }
3916 }
3917 }
3918 }
3919 if moved == 0 {
3920 continue;
3921 }
3922 let already = p
3923 .outcomes
3924 .get(&rec_hash)
3925 .and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
3926 .is_some_and(|o| o.current == moved as f64);
3927 if !already {
3928 p.outcomes.entry(rec_hash.clone()).or_default().push(
3929 crate::recommendation::OutcomeResult {
3930 rec_hash: rec_hash.clone(),
3931 metric: PREMISE_DRIFT_METRIC.into(),
3932 baseline: 0.0,
3933 current: moved as f64,
3934 verdict: "drifted".into(),
3935 baseline_kind: "snapshot".into(),
3936 baseline_run_id: None,
3937 best_before: None,
3938 tolerance: 0.0,
3939 current_run_id: None,
3940 cost: None,
3941 horizon_ms: 0,
3942 checkpoint: None,
3943 measured_at_ms: now_ms,
3944 },
3945 );
3946 }
3947 out.push(OutcomeInput {
3948 rec_hash,
3949 target_ref,
3950 metric: PREMISE_DRIFT_METRIC.into(),
3951 baseline: 0.0,
3952 current: moved as f64,
3953 unit: "superseded premises".into(),
3954 higher_is_better: false,
3955 baseline_kind: "snapshot".into(),
3956 baseline_run_id: None,
3957 best_before: None,
3958 tolerance: 0.0,
3959 current_run_id: None,
3960 cost: None,
3961 });
3962 }
3963 Ok(out)
3964}
3965
3966fn moved_premises<S: OmsSubstrate>(
3971 sub: &S,
3972 rec: &Recommendation,
3973 own: &[String],
3974) -> Result<(u64, u64)> {
3975 let total = rec.evidence.len() as u64;
3976 let mut moved = 0u64;
3977 for e in &rec.evidence {
3978 match sub.grain(e)? {
3979 None => moved += 1, Some(g) => {
3981 let Some(newer) = &g.superseded_by else { continue };
3982 if own.iter().any(|c| c == newer) {
3983 continue;
3984 }
3985 match sub.grain(newer)? {
3986 None => moved += 1,
3987 Some(n) => {
3988 if !same_value(&g, &n) {
3989 moved += 1;
3990 }
3991 }
3992 }
3993 }
3994 }
3995 }
3996 Ok((moved, total))
3997}
3998
3999fn withdraw_drifted_open<S: OmsSubstrate>(
4017 sub: &mut S,
4018 p: &mut LoopPersisted,
4019 require_all: bool,
4020 now_ms: i64,
4021) -> Result<u64> {
4022 let open: Vec<String> = p
4023 .status_index
4024 .iter()
4025 .filter(|(_, st)| matches!(st, RecStatus::Pending | RecStatus::Approved))
4026 .map(|(h, _)| h.clone())
4027 .collect();
4028 let mut withdrawn = 0u64;
4029 for rec_hash in open {
4030 let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
4031 if rec.evidence.is_empty() {
4032 continue;
4033 }
4034 let (moved, total) = moved_premises(sub, &rec, &[])?;
4035 let enough = if require_all { moved >= total } else { moved > 0 };
4036 if moved == 0 || !enough {
4037 continue;
4038 }
4039 let from = p.status_index.get(&rec_hash).copied().unwrap_or(RecStatus::Pending);
4040 let prev = p.audit_heads.get(&rec_hash).cloned();
4041 let audit = AuditRecord {
4042 rec_hash: rec_hash.clone(),
4043 from: Some(from),
4044 to: RecStatus::Withdrawn,
4045 actor: "engine:loop.premise_drift".into(),
4046 observer_type: ObserverType::System,
4047 because: format!(
4048 "{moved} of {total} cited grains were superseded by a different value or retracted"
4049 ),
4050 previous_audit_hash: prev,
4051 gating: None,
4052 at_ms: now_ms,
4053 };
4054 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
4055 p.audit_heads.insert(rec_hash.clone(), audit_hash);
4056 p.status_index.insert(rec_hash, RecStatus::Withdrawn);
4057 withdrawn += 1;
4058 }
4059 Ok(withdrawn)
4060}
4061
4062fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
4067 if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
4068 return normalize_ident(a) == normalize_ident(b);
4069 }
4070 for key in ["content", "tool_content", "body", "text", "object"] {
4071 if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
4072 return normalize_ident(a) == normalize_ident(b);
4073 }
4074 }
4075 false
4076}
4077
4078fn sanitize_skill_name(s: &str) -> Option<String> {
4080 let t = s.trim();
4081 if t.is_empty()
4082 || t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
4083 || !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
4084 {
4085 return None;
4086 }
4087 Some(t.to_string())
4088}
4089
4090struct SkillFields {
4094 fields: serde_json::Map<String, Value>,
4095 name: String,
4096 n_steps: usize,
4097 existing: Option<String>,
4098}
4099
4100#[allow(clippy::too_many_arguments)]
4105fn derived_skill_fields<S: SubstrateRead>(
4106 sub: &S,
4107 target: &TargetRef,
4108 description: &str,
4109 when_to_use: &str,
4110 steps: &[String],
4111 cited: &[String],
4112 ns_by_hash: &std::collections::BTreeMap<String, String>,
4113 skills: &crate::policy::SkillAuthoring,
4114) -> Option<SkillFields> {
4115 if target.scheme() != "entity" {
4116 return None;
4117 }
4118 let name = sanitize_skill_name(
4119 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
4120 )?;
4121 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
4122 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
4123 let steps: Vec<String> = steps
4124 .iter()
4125 .map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
4126 .filter(|st| !st.is_empty())
4127 .take(crate::llm::MAX_SKILL_STEPS)
4128 .collect();
4129 if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
4130 return None;
4131 }
4132 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
4135 for h in cited {
4136 if let Some(ns) = ns_by_hash.get(h) {
4137 if !ns.is_empty() {
4138 *ns_counts.entry(ns.as_str()).or_default() += 1;
4139 }
4140 }
4141 }
4142 let ns = ns_counts
4143 .iter()
4144 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
4145 .map(|(ns, _)| ns.to_string());
4146 let instructions = steps
4147 .iter()
4148 .enumerate()
4149 .map(|(i, st)| format!("{}. {st}", i + 1))
4150 .collect::<Vec<_>>()
4151 .join("\n");
4152 let mut fields = serde_json::Map::new();
4153 fields.insert("name".into(), Value::from(name.clone()));
4154 fields.insert("description".into(), Value::from(description));
4155 fields.insert("when_to_use".into(), Value::from(when_to_use));
4156 fields.insert("instructions".into(), Value::from(instructions));
4157 if let Some(ns) = &ns {
4158 fields.insert("namespace".into(), Value::from(ns.clone()));
4159 }
4160 let existing = sub
4162 .grains_of_type(
4163 crate::model::grain_type::SKILL,
4164 ns.as_deref(),
4165 ReadOpts { live_only: true, since_ms: None },
4166 )
4167 .ok()?
4168 .into_iter()
4169 .find(|g| g.skill_name() == Some(name.as_str()))
4170 .map(|g| g.hash);
4171 Some(SkillFields { fields, name, n_steps: steps.len(), existing })
4172}
4173
4174fn derived_fact_fields(
4175 target: &TargetRef,
4176 relation: &str,
4177 object: &str,
4178 cited: &[String],
4179 ns_by_hash: &std::collections::BTreeMap<String, String>,
4180) -> Option<serde_json::Map<String, Value>> {
4181 if target.scheme() != "entity" {
4182 return None;
4183 }
4184 let subject = target
4185 .opaque()
4186 .rsplit_once('/')
4187 .map(|(_, s)| s)
4188 .unwrap_or(target.opaque());
4189 if subject.is_empty() {
4190 return None;
4191 }
4192 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
4193 for h in cited {
4194 if let Some(ns) = ns_by_hash.get(h) {
4195 if !ns.is_empty() {
4196 *ns_counts.entry(ns.as_str()).or_default() += 1;
4197 }
4198 }
4199 }
4200 let lesson_ns = ns_counts
4201 .iter()
4202 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
4203 .map(|(ns, _)| ns.to_string());
4204 let mut fields = serde_json::Map::new();
4205 fields.insert("subject".into(), Value::from(subject));
4206 fields.insert("relation".into(), Value::from(relation));
4207 fields.insert("object".into(), Value::from(object));
4208 if let Some(ns) = lesson_ns {
4212 fields.insert("namespace".into(), Value::from(ns));
4213 }
4214 Some(fields)
4215}
4216
4217fn latest_verdicts(p: &LoopPersisted) -> BTreeMap<String, String> {
4222 let mut out = BTreeMap::new();
4223 for (rec_hash, applied) in &p.applied {
4224 let latest = p
4225 .outcomes
4226 .get(rec_hash)
4227 .and_then(|v| v.iter().max_by_key(|o| o.measured_at_ms))
4228 .map(|o| o.verdict.clone());
4229 for h in &applied.created_hashes {
4230 out.insert(h.clone(), latest.clone().unwrap_or_else(|| "unmeasured".into()));
4231 }
4232 }
4233 out
4234}
4235
4236pub const NEAR_DUPLICATE_COSINE: f64 = 0.90;
4239pub const NEAR_DUPLICATE_JACCARD: f64 = 0.60;
4242const NEAR_DUPLICATE_CAP: usize = 8;
4244
4245pub(crate) fn near_duplicates_of<S: SubstrateRead + ?Sized>(
4250 sub: &S,
4251 subject: &str,
4252 namespace: Option<&str>,
4253 text: &str,
4254) -> Vec<crate::recommendation::NearDuplicate> {
4255 use crate::analyzers::duplicate_sweep::{jaccard, tokenize};
4256 if subject.is_empty() || text.trim().is_empty() {
4257 return Vec::new();
4258 }
4259 let Ok(facts) = sub.grains_of_type(
4260 crate::model::grain_type::FACT,
4261 namespace,
4262 ReadOpts { live_only: true, since_ms: None },
4263 ) else {
4264 return Vec::new();
4265 };
4266 let mine = sub.embed(text).ok().flatten();
4267 let my_tokens = tokenize(text);
4268 let mut out: Vec<crate::recommendation::NearDuplicate> = facts
4269 .iter()
4270 .filter(|f| f.fact_relation() == Some("lesson"))
4271 .filter(|f| f.fact_subject().is_some_and(|s| normalize_ident(s) == normalize_ident(subject)))
4272 .filter_map(|f| {
4273 let other = f.fact_object()?;
4274 let (score, method, floor) = match (&mine, sub.embed(other).ok().flatten()) {
4275 (Some(a), Some(b)) => (cosine(a, &b), "cosine", NEAR_DUPLICATE_COSINE),
4276 _ => (jaccard(&my_tokens, &tokenize(other)), "jaccard", NEAR_DUPLICATE_JACCARD),
4277 };
4278 (score >= floor).then(|| crate::recommendation::NearDuplicate {
4279 hash: f.hash.clone(),
4280 score: (score * 1000.0).round() / 1000.0,
4281 method: method.into(),
4282 })
4283 })
4284 .collect();
4285 out.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal).then(a.hash.cmp(&b.hash)));
4286 out.truncate(NEAR_DUPLICATE_CAP);
4287 out
4288}
4289
4290fn cosine(a: &[f32], b: &[f32]) -> f64 {
4291 if a.len() != b.len() || a.is_empty() {
4292 return 0.0;
4293 }
4294 let (mut dot, mut na, mut nb) = (0f64, 0f64, 0f64);
4295 for (x, y) in a.iter().zip(b) {
4296 dot += *x as f64 * *y as f64;
4297 na += *x as f64 * *x as f64;
4298 nb += *y as f64 * *y as f64;
4299 }
4300 if na == 0.0 || nb == 0.0 {
4301 0.0
4302 } else {
4303 dot / (na.sqrt() * nb.sqrt())
4304 }
4305}
4306
4307const CONSOLIDATION_INSTRUCTIONS: &str = " (8) {\"kind\":\"consolidation\",\"lesson\":\"...\",\
4311\"supersedes\":[\"<hash>\",...]} with the same entity target — ONLY in answer to a \
4312'Lesson pile' finding, which lists the live lessons on one entity that exceed \
4313its budget. Write ONE short imperative rule (max 240 chars) that says what \
4314those lessons say together, dropping nothing a lesson that measured 'held' \
4315required and keeping nothing only a lesson that measured 'regressed' or \
4316'drifted' added; 'supersedes' MUST be exactly the hashes that finding lists \
4317(cite them as evidence too). Applying it replaces every listed lesson with \
4318the one line; the reviewer can restore them all.";
4319
4320fn requires_gating(kind: ActionKind) -> bool {
4325 matches!(
4326 kind,
4327 ActionKind::CodeRevision | ActionKind::AdapterRevision
4328 )
4329}
4330
4331fn stamp(
4332 m: &AnalyzerManifest,
4333 params: &crate::manifest::Params,
4334 d: crate::recommendation::RecDraft,
4335 now_ms: i64,
4336 scope: &[String],
4337) -> Result<Recommendation> {
4338 let target = TargetRef::parse(&d.target_ref)?;
4339 crate::recommendation::validate_code_rules(
4343 d.action_kind,
4344 target.target_class(),
4345 d.evalset_hash.as_deref(),
4346 )?;
4347 let revert_of = match (&d.action_kind, &d.proposal) {
4350 (ActionKind::Revert, Proposal::Data { data }) => {
4351 data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
4352 }
4353 _ => None,
4354 };
4355 let dedup = match revert_of.as_deref() {
4356 Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
4357 None => dedup_key(m.family(), &d.target_ref, d.action_kind),
4358 };
4359 let destructive = match &d.proposal {
4360 Proposal::Cal { cal } => cal::contains_destructive(cal),
4361 _ => false,
4362 };
4363 let rollbackable = match &d.proposal {
4364 Proposal::Cal { .. } => !destructive,
4365 Proposal::Edit { .. } => false,
4366 Proposal::Data { .. } => requires_gating(d.action_kind),
4370 };
4371 let mut evidence = d.evidence;
4372 evidence.truncate(MAX_EVIDENCE);
4373 let origin = match m.trust_class {
4378 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
4379 _ => Origin::Builtin,
4380 };
4381 Ok(Recommendation {
4382 hash: String::new(),
4383 analyzer: m.id.clone(),
4384 params_snapshot: params.snapshot(),
4385 origin,
4386 target_ref: target.as_string(),
4387 action_kind: d.action_kind,
4388 dedup_key: dedup,
4389 summary: d.summary,
4390 severity: d.severity,
4391 proposal: d.proposal,
4392 destructive,
4393 rollbackable,
4394 evidence,
4395 evidence_query: d.evidence_query,
4396 metric: d.metric,
4397 confidence: d.confidence,
4398 importance: d.importance,
4399 created_at_ms: now_ms,
4400 guidance: None,
4401 evalset_hash: d.evalset_hash,
4402 near_duplicate_of: Vec::new(),
4403 replay: None,
4404 scope: crate::recommendation::normalize_scope(scope),
4409 status: RecStatus::Pending,
4410 })
4411}
4412
4413fn validate_because(because: &str) -> Result<String> {
4414 let trimmed = because.trim();
4415 if trimmed.is_empty() {
4416 return Err(Error::InvalidProposal(
4417 "a BECAUSE reason is required".into(),
4418 ));
4419 }
4420 if trimmed.chars().count() > MAX_BECAUSE {
4421 return Err(Error::InvalidProposal(format!(
4422 "BECAUSE exceeds {MAX_BECAUSE} chars"
4423 )));
4424 }
4425 Ok(trimmed.to_string())
4426}
4427
4428fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
4429 for req in &m.requires {
4430 match req {
4431 Capability::Forks if !caps.forks => return Some("forks"),
4432 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
4433 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
4434 _ => {}
4435 }
4436 }
4437 None
4438}
4439
4440fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
4441 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
4442}
4443
4444fn gate(
4445 opts: &RunOptions,
4446 p: &LoopPersisted,
4447 new_grains: u64,
4448 new_errors: u64,
4449 now_ms: i64,
4450) -> Option<SkipReason> {
4451 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
4452 if !any {
4453 return None;
4454 }
4455 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
4456 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
4457 let stale_ok = opts
4458 .if_stale_ms
4459 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
4460 if min_new_ok || min_err_ok || stale_ok {
4461 return None;
4462 }
4463 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
4465 Some(SkipReason::NotStale)
4466 } else {
4467 Some(SkipReason::MinNewNotMet)
4468 }
4469}
4470
4471#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
4473pub(crate) struct NewSince {
4474 pub grains: u64,
4476 pub error_events: u64,
4478 pub events: u64,
4480 pub sessions: u64,
4482}
4483
4484fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
4485 let opts = ReadOpts {
4486 live_only: false,
4487 since_ms: watermark.map(|w| w + 1),
4488 };
4489 let mut n = NewSince::default();
4490 let mut sessions: BTreeSet<&str> = BTreeSet::new();
4491 let mut events_held: Vec<GrainRecord> = Vec::new();
4492 for t in [
4493 crate::model::grain_type::FACT,
4494 crate::model::grain_type::EVENT,
4495 crate::model::grain_type::TOOL,
4496 crate::model::grain_type::OBSERVATION,
4497 ] {
4498 let g = sub.grains_of_type(t, None, opts)?;
4499 n.grains += g.len() as u64;
4500 if t == crate::model::grain_type::TOOL {
4502 n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
4503 }
4504 if t == crate::model::grain_type::EVENT {
4505 n.events = g.len() as u64;
4506 events_held = g;
4507 }
4508 }
4509 for e in &events_held {
4510 if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
4511 sessions.insert(sid);
4512 }
4513 }
4514 n.sessions = sessions.len() as u64;
4515 Ok(n)
4516}
4517
4518fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
4525 if !c.is_set() {
4526 return None;
4527 }
4528 let time_ok = c
4529 .every_ms
4530 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
4531 let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
4532 let events_ok = c.every_events.is_some_and(|m| new.events >= m);
4533 let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
4534 if time_ok || grains_ok || events_ok || sessions_ok {
4535 None
4536 } else {
4537 Some(SkipReason::CadenceNotDue)
4538 }
4539}
4540
4541fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
4542 let grains = sub.grains_of_type(
4543 crate::model::grain_type::RECOMMENDATION,
4544 Some(LOOP_NS),
4545 ReadOpts {
4546 live_only: false,
4547 since_ms: None,
4548 },
4549 )?;
4550 let mut set = BTreeSet::new();
4551 for g in grains {
4552 let status = p
4553 .status_index
4554 .get(&g.hash)
4555 .copied()
4556 .unwrap_or(RecStatus::Pending);
4557 if matches!(
4571 status,
4572 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
4573 ) {
4574 if let Some(key) = g.str_field("dedup_key") {
4575 set.insert(key.to_string());
4576 }
4577 }
4578 }
4579 Ok(set)
4580}
4581
4582pub(crate) fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
4589 const BASE_MS: i64 = 7 * 86_400_000;
4590 const CAP_MS: i64 = 90 * 86_400_000;
4591 let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
4592 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
4593 *strikes = strikes.saturating_add(1);
4594 p.cooldowns.insert(dedup_key, now_ms + interval);
4595}
4596
4597fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
4598 let g = sub
4599 .grain(rec_hash)?
4600 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
4601 Recommendation::from_fields(rec_hash, &g.fields)
4602}
4603
4604pub(crate) fn is_definition_statement(line: &str) -> bool {
4612 let up = line.trim_start().to_ascii_uppercase();
4613 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
4614}
4615
4616const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
4619 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
4620 approve it to acknowledge it and let it expire.";
4621
4622const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
4627 (evalset hash + run id + stats) — use apply_gated";
4628
4629const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
4630 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
4631 acknowledge it and let it expire.";
4632
4633pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
4646 match proposal {
4647 Proposal::Cal { .. } => Ok(()),
4648 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
4649 Proposal::Data { data } => {
4655 if requires_gating(action_kind)
4656 || data.get("revert_of").and_then(Value::as_str).is_some()
4657 {
4658 Ok(())
4659 } else {
4660 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
4661 }
4662 }
4663 }
4664}
4665
4666#[cfg(test)]
4667mod definition_body_tests {
4668 use super::safe_definition_body;
4669
4670 #[test]
4671 fn ordinary_bodies_pass() {
4672 assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
4673 assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
4674 }
4675
4676 #[test]
4677 fn a_body_cannot_close_its_own_block_or_carry_destruction() {
4678 assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
4683 assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
4684 assert!(!safe_definition_body("RECALL facts FORGET abc"));
4686 assert!(!safe_definition_body("recall facts purge older than 1d"));
4687 assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
4688 assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
4690 assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
4692 }
4693}
4694
4695#[cfg(test)]
4696mod plan_edit_tests {
4697 use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
4698 use serde_json::json;
4699
4700 fn plan() -> serde_json::Value {
4701 json!({
4702 "nodes": ["fetch", "review", "post"],
4703 "edges": [
4704 {"src": "fetch", "dst": "review"},
4705 {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
4706 ],
4707 "bindings": {"fetch": "sha256:tool1"},
4708 "retries": {"fetch": 1}
4709 })
4710 }
4711
4712 #[test]
4713 fn the_allowlist_admits_thresholds_and_refuses_topology() {
4714 assert!(plan_edit_allowed("edges.1.cond"));
4715 assert!(plan_edit_allowed("edges.1.max_cycles"));
4716 assert!(plan_edit_allowed("retries.fetch"));
4717 for path in [
4720 "nodes",
4721 "nodes.0",
4722 "edges.0.src",
4723 "edges.0.dst",
4724 "edges",
4725 "bindings.fetch",
4726 "edges.x.cond",
4727 "",
4728 ] {
4729 assert!(!plan_edit_allowed(path), "{path} must not be editable");
4730 }
4731 }
4732
4733 #[test]
4734 fn values_are_type_checked_against_the_field() {
4735 assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
4738 assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
4739 assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
4740 assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
4741 assert!(plan_value_ok("retries.fetch", &json!(3)));
4742 assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
4743 assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
4744 assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
4745 assert!(!plan_value_ok("edges.0.src", &json!("other")));
4746 }
4747
4748 #[test]
4749 fn get_reads_through_arrays_and_objects_and_absence_is_null() {
4750 let p = plan();
4751 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
4752 assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
4753 assert_eq!(plan_get(&p, "retries.review"), json!(null));
4756 assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
4757 assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
4758 }
4759
4760 #[test]
4761 fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
4762 let mut p = plan();
4763 assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
4764 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
4765 assert!(plan_set(&mut p, "retries.review", json!(2)));
4766 assert_eq!(plan_get(&p, "retries.review"), json!(2));
4767 assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
4768 assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
4769 let mut q = plan();
4774 q.as_object_mut().unwrap().remove("retries");
4775 assert!(plan_set(&mut q, "retries.greet", json!(1)));
4776 assert_eq!(plan_get(&q, "retries.greet"), json!(1));
4777 let mut r = plan();
4779 r.as_object_mut().unwrap().remove("edges");
4780 assert!(!plan_set(&mut r, "edges.0.cond", json!("x")));
4781 }
4782}
4783
4784#[cfg(test)]
4785mod definition_proposal_tests {
4786 use super::is_definition_statement;
4787
4788 #[test]
4789 fn definition_statements_are_recognized_in_both_spellings() {
4790 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
4791 assert!(is_definition_statement(" define template foo AS { x }"));
4792 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
4793 assert!(!is_definition_statement("ADD fact {}"));
4795 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
4796 assert!(!is_definition_statement("FORGET abc"));
4797 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
4801 }
4802}