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}
136
137#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
148pub struct LlmFunnel {
149 pub evidence: u64,
151 pub proposed: u64,
153 pub cited: u64,
155 pub dropped_uncited: u64,
159 pub dropped_target: u64,
162 pub grounded: u64,
164 pub ground_verdicts: u64,
169 pub ground_call_failed: bool,
173 pub kept: u64,
175 pub stored: u64,
177 #[serde(default, skip_serializing_if = "is_zero")]
181 pub advisory_thin_evidence: u64,
182}
183
184fn is_zero(n: &u64) -> bool {
185 *n == 0
186}
187
188impl RunResult {
189 fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
190 RunResult {
191 outcome: RunOutcome::Skipped,
192 skip_reason: Some(reason),
193 new_grains,
194 new_error_events,
195 proposed: 0,
196 deduped: 0,
197 stored: 0,
198 auto_applied: 0,
199 llm_funnel: None,
200 analyzers_run: vec![],
201 analyzers_skipped: vec![],
202 }
203 }
204
205 pub fn ran(&self) -> bool {
206 self.outcome == RunOutcome::Ran
207 }
208}
209
210pub struct Engine {
213 analyzers: Vec<Box<dyn Analyzer>>,
214 policy: crate::policy::Policy,
215 llm: Option<Box<dyn crate::llm::LlmBackend>>,
218 ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
223}
224
225struct AnalysisPass {
226 survivors: Vec<Recommendation>,
227 proposed: u64,
228 deduped: u64,
229 analyzers_run: Vec<String>,
230 analyzers_skipped: Vec<AnalyzerSkip>,
231 llm_funnel: Option<LlmFunnel>,
232}
233
234impl Engine {
235 pub fn with_builtins() -> Self {
238 Engine {
239 analyzers: crate::analyzer::builtin_analyzers(),
240 policy: crate::policy::Policy::default(),
241 llm: None,
242 ground_llm: None,
243 }
244 }
245
246 pub fn empty() -> Self {
248 Engine {
249 analyzers: vec![],
250 policy: crate::policy::Policy::default(),
251 llm: None,
252 ground_llm: None,
253 }
254 }
255
256 pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
258 self.policy = policy;
259 self
260 }
261
262 pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
267 self.llm = Some(backend);
268 self
269 }
270
271 pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
275 self.ground_llm = Some(backend);
276 self
277 }
278
279 pub fn policy(&self) -> &crate::policy::Policy {
280 &self.policy
281 }
282
283 pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
285 self.analyzers.push(analyzer);
286 }
287
288 pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
289 &self.analyzers
290 }
291
292 pub fn analyze_only<S: OmsSubstrate>(
301 &self,
302 sub: &S,
303 opts: &RunOptions,
304 overrides: &BTreeMap<String, Map<String, Value>>,
305 now_ms: i64,
306 ) -> Result<Vec<Recommendation>> {
307 let persisted = LoopPersisted::from_value(sub.load_state()?)?;
308 let analysis_watermark = if opts.full_sweep {
309 None
310 } else {
311 persisted.state.watermark_ms
312 };
313 Ok(self
314 .analysis_pass(
315 sub,
316 &persisted,
317 opts,
318 overrides,
319 analysis_watermark,
320 now_ms,
321 &[],
322 )?
323 .survivors)
324 }
325
326 pub fn run<S: OmsSubstrate>(
329 &self,
330 sub: &mut S,
331 opts: &RunOptions,
332 now_ms: i64,
333 ) -> Result<RunResult> {
334 let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
335 let watermark = persisted.state.watermark_ms;
336 let analysis_watermark = if opts.full_sweep { None } else { watermark };
341
342 let new = count_new(sub, watermark)?;
343 let (new_grains, new_error_events) = (new.grains, new.error_events);
344 if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
345 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
346 }
347 let flags_set = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
351 if !flags_set && !opts.full_sweep {
352 if let Some(reason) = cadence_gate(&self.policy.cadence, &persisted, new, now_ms) {
353 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
354 }
355 }
356
357 let mut outcome_inputs = measure_outcomes(sub, &mut persisted, now_ms)?;
362 if self.policy.premise_drift {
363 outcome_inputs.extend(detect_premise_drift(sub, &mut persisted, now_ms)?);
364 }
365
366 let AnalysisPass {
367 survivors,
368 proposed,
369 deduped,
370 analyzers_run,
371 analyzers_skipped,
372 llm_funnel,
373 } = self.analysis_pass(
374 &*sub,
375 &persisted,
376 opts,
377 &BTreeMap::new(),
378 analysis_watermark,
379 now_ms,
380 &outcome_inputs,
381 )?;
382
383 let mut stored = 0u64;
386 let mut auto_applied = 0u64;
387 for mut rec in survivors {
388 let spec = rec.to_grain_spec(LOOP_NS)?;
389 let hash = sub.put_grain(&spec)?;
390 rec.hash = hash.clone();
391 let actor = format!("engine:{}", rec.analyzer);
392 let audit = AuditRecord {
393 rec_hash: hash.clone(),
394 from: None,
395 to: RecStatus::Pending,
396 actor: actor.clone(),
397 observer_type: ObserverType::System,
398 because: "analyzer proposed".into(),
399 previous_audit_hash: None,
400 gating: None,
401 at_ms: now_ms,
402 };
403 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
404 persisted
405 .status_index
406 .insert(hash.clone(), RecStatus::Pending);
407 persisted.creators.insert(hash.clone(), actor);
408 if !matches!(rec.origin, Origin::Builtin) {
413 if let Some(trigger) = &opts.triggering_actor {
414 persisted.co_creators.insert(hash.clone(), trigger.clone());
415 }
416 }
417 persisted.audit_heads.insert(hash.clone(), audit_hash);
418 stored += 1;
419
420 if self.can_auto_apply(&*sub, &rec) {
421 self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
422 auto_applied += 1;
423 }
424 }
425
426 persisted.state.last_run_ms = Some(now_ms);
427 persisted.state.watermark_ms = Some(now_ms);
428 sub.store_state(&persisted.to_value()?)?;
429
430 Ok(RunResult {
431 outcome: RunOutcome::Ran,
432 skip_reason: None,
433 new_grains,
434 new_error_events,
435 proposed,
436 deduped,
437 stored,
438 auto_applied,
439 analyzers_run,
440 analyzers_skipped,
441 llm_funnel,
442 })
443 }
444
445 #[allow(clippy::too_many_arguments)]
449 fn analysis_pass<S: OmsSubstrate>(
450 &self,
451 sub: &S,
452 persisted: &LoopPersisted,
453 opts: &RunOptions,
454 external_overrides: &BTreeMap<String, Map<String, Value>>,
455 analysis_watermark: Option<i64>,
456 now_ms: i64,
457 outcome_inputs: &[OutcomeInput],
458 ) -> Result<AnalysisPass> {
459 let existing = existing_dedup_keys(sub, persisted)?;
460 let mut analyzers_run = Vec::new();
461 let mut analyzers_skipped = Vec::new();
462 let mut candidates: Vec<Recommendation> = Vec::new();
463 let caps = sub.capabilities();
464
465 for analyzer in &self.analyzers {
466 let m = analyzer.manifest();
467 let cfg = persisted.config.get(&m.id);
468 let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
469 if !enabled {
470 analyzers_skipped.push(AnalyzerSkip {
471 id: m.id.clone(),
472 reason: "disabled".into(),
473 });
474 continue;
475 }
476 if self.policy.denies(m.family()) {
477 analyzers_skipped.push(AnalyzerSkip {
478 id: m.id.clone(),
479 reason: "denied by host policy".into(),
480 });
481 continue;
482 }
483 if let Some(missing) = missing_capability(m, caps) {
484 analyzers_skipped.push(AnalyzerSkip {
485 id: m.id.clone(),
486 reason: format!("missing capability: {missing}"),
487 });
488 continue;
489 }
490 let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
491 if let Some(extra) = external_overrides.get(&m.id) {
492 for (key, value) in extra {
493 param_overrides.insert(key.clone(), value.clone());
494 }
495 }
496 let params = match m.resolve_params(¶m_overrides) {
497 Ok(p) => p,
498 Err(e) => {
499 analyzers_skipped.push(AnalyzerSkip {
500 id: m.id.clone(),
501 reason: e.to_string(),
502 });
503 continue;
504 }
505 };
506 let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
507 let ns_slice: &[String] = if ns_owned.is_empty() {
508 &opts.namespaces
509 } else {
510 &ns_owned
511 };
512 let reader: &dyn SubstrateRead = sub;
513 let ctx = AnalyzeCtx::new(
514 reader,
515 ¶ms,
516 ns_slice,
517 analysis_watermark,
518 now_ms,
519 outcome_inputs,
520 );
521 match analyzer.analyze(&ctx) {
522 Ok(drafts) => {
523 analyzers_run.push(m.id.clone());
524 for draft in drafts {
525 match stamp(m, ¶ms, draft, now_ms) {
526 Ok(rec) => candidates.push(rec),
527 Err(e) => analyzers_skipped.push(AnalyzerSkip {
528 id: m.id.clone(),
529 reason: e.to_string(),
530 }),
531 }
532 }
533 }
534 Err(e) => analyzers_skipped.push(AnalyzerSkip {
535 id: m.id.clone(),
536 reason: e.to_string(),
537 }),
538 }
539 }
540
541 let mut funnel = LlmFunnel::default();
542 if self.llm.is_some() {
543 candidates.extend(self.discover(
544 sub,
545 &candidates,
546 analysis_watermark,
547 &opts.namespaces,
548 now_ms,
549 &mut funnel,
550 ));
551 }
552
553 let proposed = candidates.len() as u64;
554 let mut seen = BTreeSet::new();
555 let mut survivors = Vec::new();
556 for candidate in candidates {
557 let family = crate::manifest::analyzer_family(&candidate.analyzer);
558 let floor = [
559 severity_floor_for(persisted, &candidate.analyzer),
560 self.policy.severity_floor(family),
561 ]
562 .into_iter()
563 .flatten()
564 .max();
565 if floor.is_some_and(|floor| candidate.severity < floor) {
566 continue;
567 }
568 if !seen.insert(candidate.dedup_key.clone()) {
569 continue;
570 }
571 if existing.contains(&candidate.dedup_key) {
572 continue;
573 }
574 if persisted
575 .cooldowns
576 .get(&candidate.dedup_key)
577 .is_some_and(|until| now_ms < *until)
578 {
579 continue;
580 }
581 survivors.push(candidate);
582 }
583 let deduped = proposed - survivors.len() as u64;
584 if self.llm.is_some() {
585 self.enrich(&mut survivors);
586 }
587 Ok(AnalysisPass {
588 survivors,
589 proposed,
590 deduped,
591 analyzers_run,
592 analyzers_skipped,
593 llm_funnel: self.llm.is_some().then_some(funnel),
594 })
595 }
596
597 fn discover<S: OmsSubstrate>(
604 &self,
605 sub: &S,
606 candidates: &[Recommendation],
607 watermark: Option<i64>,
608 namespaces: &[String],
609 now_ms: i64,
610 funnel: &mut LlmFunnel,
611 ) -> Vec<Recommendation> {
612 let Some(llm) = &self.llm else {
613 return Vec::new();
614 };
615 let findings: Vec<crate::llm::FindingBrief> = candidates
616 .iter()
617 .take(32)
618 .map(|c| crate::llm::FindingBrief {
619 analyzer: c.analyzer.clone(),
620 summary: c.summary.render(),
621 target: c.target_ref.clone(),
622 severity: c.severity.as_str().to_string(),
623 })
624 .collect();
625 let attribution = self.policy.evidence_attribution;
631 let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
632 let mut bundle: BTreeSet<String> = BTreeSet::new();
633 let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
634 'cited: for c in candidates {
635 for h in &c.evidence {
636 if evidence.len() >= CITED_SEED_CAP {
645 break 'cited;
646 }
647 if !bundle.contains(h) {
648 if let Ok(Some(g)) = sub.grain(h) {
649 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
650 }
651 }
652 }
653 }
654 let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
655 vec![None]
656 } else {
657 namespaces.iter().map(|n| Some(n.as_str())).collect()
658 };
659 let opts = ReadOpts { live_only: true, since_ms: watermark };
660 let mut tool_seeded = 0usize;
681 let seed_successes = self.policy.skills.enabled || self.policy.plans.enabled;
682 'tools: for want_error in [true, false] {
683 if !want_error && !seed_successes {
684 break;
685 }
686 for ns in &scan_ns {
687 if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
688 for g in recent {
689 if tool_seeded >= TOOL_SEED_CAP
690 || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
691 {
692 break 'tools;
693 }
694 if g.is_error() != want_error {
695 continue;
696 }
697 let before = evidence.len();
698 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
699 if evidence.len() > before {
700 tool_seeded += 1;
701 }
702 }
703 }
704 }
705 }
706 'notes: for ns in &scan_ns {
716 if let Ok(recent) =
717 sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
718 {
719 for g in recent {
720 if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
721 break 'notes;
722 }
723 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
724 }
725 }
726 }
727 'seed: for gt in [
728 crate::model::grain_type::FACT,
729 crate::model::grain_type::OBSERVATION,
730 ] {
731 for ns in &scan_ns {
732 if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
733 for g in recent {
734 if evidence.len() >= EVIDENCE_CAP {
735 break 'seed;
736 }
737 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
738 }
739 }
740 }
741 }
742 funnel.evidence = evidence.len() as u64;
743 if evidence.is_empty() {
744 return Vec::new(); }
746 let (approved, rejected) = self.llm_history(sub);
751 let base = match self.policy.discover_objective {
752 crate::policy::DiscoverObjective::ReviewQueue => DISCOVER_INSTRUCTIONS,
753 crate::policy::DiscoverObjective::Learner => DISCOVER_LEARNER_INSTRUCTIONS,
754 };
755 let mut instructions = base.to_string();
758 if self.policy.skills.enabled {
759 instructions.push_str(&skill_instructions(self.policy.skills.min_steps));
760 }
761 if self.policy.plans.enabled {
762 instructions.push_str(&plan_instructions(self.policy.plans.min_nodes));
763 }
764 let request = crate::llm::LlmRequest {
765 loop_proto: 1,
766 op: "discover",
767 instructions: &instructions,
768 findings: findings.clone(),
769 evidence: evidence.clone(),
770 rejected,
771 approved,
772 };
773 let Ok(body) = serde_json::to_string(&request) else {
774 return Vec::new();
775 };
776 let raw = match llm.complete(&body) {
777 Ok(r) => r,
778 Err(_) => return Vec::new(), };
780 let caps = sub.capabilities();
784 let mut validated: Vec<ValidatedDraft> = Vec::new();
785 let drafts: Vec<_> = crate::llm::parse_discover(&raw)
786 .recommendations
787 .into_iter()
788 .take(crate::llm::MAX_LLM_DRAFTS)
789 .collect();
790 funnel.proposed = drafts.len() as u64;
791 let id_to_hash: std::collections::BTreeMap<&str, &str> = evidence
792 .iter()
793 .map(|e| (e.id.as_str(), e.hash.as_str()))
794 .collect();
795 for d in drafts {
796 let mut cited: Vec<String> = Vec::new();
797 for c in &d.evidence {
798 if let Some(h) = resolve_citation(c, &bundle, &id_to_hash) {
799 if !cited.contains(&h) {
800 cited.push(h);
801 }
802 }
803 }
804 if cited.is_empty() {
805 funnel.dropped_uncited += 1;
806 continue; }
808 let Ok(target) = TargetRef::parse(&d.target) else {
809 funnel.dropped_target += 1;
810 continue;
811 };
812 let tc = target.target_class();
813 if !matches!(tc, "memory" | "query" | "code") {
819 funnel.dropped_target += 1;
820 continue;
821 }
822 let thin = cited.len() < self.policy.min_evidence as usize;
829 if thin {
830 funnel.advisory_thin_evidence += 1;
831 }
832 let resolved = if thin {
833 None
834 } else {
835 resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps, &self.policy)
836 };
837 if tc == "code" && resolved.is_none() {
842 funnel.dropped_target += 1;
843 continue;
844 }
845 validated.push(ValidatedDraft {
846 draft: d,
847 target_ref: target.as_string(),
848 cited,
849 resolved,
850 });
851 }
852 funnel.cited = validated.len() as u64;
853 if validated.is_empty() {
854 return Vec::new();
855 }
856 let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
863 let outcome_metric = self.outcome_metric_template(sub);
864 self.verify_drafts(&**llm, ground, validated, &evidence, outcome_metric, now_ms, funnel)
865 }
866
867 fn outcome_metric_template<S: OmsSubstrate>(
875 &self,
876 sub: &S,
877 ) -> Option<crate::recommendation::MetricSnapshot> {
878 let e = self.policy.outcome_evalset.as_ref()?;
879 let run = crate::eval::newest_eval_run(sub, &e.hash, None).ok().flatten()?;
880 let baseline = crate::eval::run_value(&run, &e.field)?;
881 let schedule = e.schedule();
885 let ms_only: Vec<i64> = schedule.iter().filter_map(Checkpoint::as_ms).collect();
886 let all_ms = ms_only.len() == schedule.len();
887 let horizons = if ms_only.is_empty() { vec![86_400_000] } else { ms_only };
888 Some(crate::recommendation::MetricSnapshot {
889 metric: format!("evalset:{}:{}", e.hash, e.field),
890 baseline,
891 unit: e.field.clone(),
892 n: run.total(),
893 window: "per-run".into(),
894 subject: None,
895 namespace: None,
896 relation: None,
897 query: format!(
898 "RECALL facts WHERE subject = \"evalset:{}\" AND relation = \"mg:eval_run\"",
899 e.hash
900 ),
901 review_after_ms: horizons[0],
902 horizons_ms: if all_ms { horizons } else { Vec::new() },
903 checkpoints: if all_ms { Vec::new() } else { schedule },
904 higher_is_better: e.higher_is_better,
905 })
906 }
907
908 #[allow(clippy::too_many_arguments)]
915 fn verify_drafts(
916 &self,
917 llm: &dyn crate::llm::LlmBackend,
918 ground: &dyn crate::llm::LlmBackend,
919 validated: Vec<ValidatedDraft>,
920 evidence: &[crate::llm::EvidenceItem],
921 outcome_metric: Option<crate::recommendation::MetricSnapshot>,
922 now_ms: i64,
923 funnel: &mut LlmFunnel,
924 ) -> Vec<Recommendation> {
925 use crate::llm::*;
926 let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
927 evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
928 let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
929 cited
930 .iter()
931 .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
932 .collect()
933 };
934
935 let claims: Vec<GroundItem> = validated
939 .iter()
940 .enumerate()
941 .map(|(i, v)| GroundItem {
942 id: i,
943 claim: claim_text(&v.draft, v.resolved.as_ref()),
944 evidence: ev_for(&v.cited),
945 })
946 .collect();
947 let ground_req = GroundRequest {
948 loop_proto: 1,
949 op: "ground",
950 instructions: GROUND_INSTRUCTIONS,
951 claims,
952 };
953 let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
959 .ok()
960 .and_then(|b| ground.complete(&b).ok())
961 {
962 Some(raw) => {
963 let parsed = parse_ground(&raw);
964 funnel.ground_verdicts = parsed.results.len() as u64;
965 parsed
966 .results
967 .into_iter()
968 .filter(|r| r.supported)
969 .map(|r| r.id)
970 .collect()
971 }
972 None => {
973 funnel.ground_call_failed = true;
974 return Vec::new();
975 }
976 };
977 funnel.grounded = grounded.len() as u64;
978 if grounded.is_empty() {
979 return Vec::new();
980 }
981
982 let items: Vec<VerifyItem> = validated
988 .iter()
989 .enumerate()
990 .filter(|(i, _)| grounded.contains(i))
991 .map(|(i, v)| VerifyItem {
992 id: i,
993 summary: claim_text(&v.draft, v.resolved.as_ref()),
995 target: v.target_ref.clone(),
996 evidence: ev_for(&v.cited),
997 })
998 .collect();
999 let verify_req = VerifyRequest {
1000 loop_proto: 1,
1001 op: "verify",
1002 instructions: VERIFY_INSTRUCTIONS,
1003 findings: items,
1004 };
1005 let verdicts: std::collections::BTreeMap<usize, f64> =
1006 match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
1007 Some(raw) => parse_verify(&raw)
1008 .results
1009 .into_iter()
1010 .filter(|r| r.keep)
1011 .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
1012 .collect(),
1013 None => return Vec::new(),
1014 };
1015
1016 funnel.kept = verdicts.len() as u64;
1017 let mut out = Vec::new();
1021 for (i, v) in validated.into_iter().enumerate() {
1022 if let Some(&conf) = verdicts.get(&i) {
1023 if conf >= MIN_LLM_CONFIDENCE {
1024 let mut rec = stamp_llm(
1025 llm.model(),
1026 &v.draft,
1027 v.target_ref,
1028 v.cited,
1029 v.resolved,
1030 conf,
1031 now_ms,
1032 );
1033 if rec.rollbackable {
1037 rec.metric = outcome_metric.clone();
1038 }
1039 out.push(rec);
1040 }
1041 }
1042 }
1043 funnel.stored = out.len() as u64;
1044 out
1045 }
1046
1047 fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
1052 const MAX: usize = 20;
1053 let Ok(mut recs) = self.recommendations(sub, None) else {
1054 return (Vec::new(), Vec::new());
1055 };
1056 recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
1057 recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
1058 let mut approved = Vec::new();
1059 let mut rejected = Vec::new();
1060 for r in &recs {
1061 match r.status {
1062 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
1063 if approved.len() < MAX =>
1064 {
1065 approved.push(r.summary.render());
1066 }
1067 RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
1068 _ => {}
1069 }
1070 }
1071 (approved, rejected)
1072 }
1073
1074 fn enrich(&self, survivors: &mut [Recommendation]) {
1079 let Some(llm) = &self.llm else {
1080 return;
1081 };
1082 if survivors.is_empty() {
1083 return;
1084 }
1085 let findings: Vec<crate::llm::FindingBrief> = survivors
1086 .iter()
1087 .map(|r| crate::llm::FindingBrief {
1088 analyzer: r.analyzer.clone(),
1089 summary: r.summary.render(),
1090 target: r.target_ref.clone(),
1091 severity: r.severity.as_str().to_string(),
1092 })
1093 .collect();
1094 let request = crate::llm::LlmRequest {
1095 loop_proto: 1,
1096 op: "enrich",
1097 instructions: ENRICH_INSTRUCTIONS,
1098 findings,
1099 evidence: Vec::new(),
1100 rejected: Vec::new(),
1101 approved: Vec::new(),
1102 };
1103 let Ok(body) = serde_json::to_string(&request) else {
1104 return;
1105 };
1106 let raw = match llm.complete(&body) {
1107 Ok(r) => r,
1108 Err(_) => return,
1109 };
1110 for note in crate::llm::parse_enrich(&raw).notes {
1111 if note.guidance.trim().is_empty() {
1112 continue;
1113 }
1114 if let Some(r) = survivors
1115 .iter_mut()
1116 .find(|r| r.target_ref == note.target && r.guidance.is_none())
1117 {
1118 r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
1119 }
1120 }
1121 }
1122
1123 fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
1132 if !rec.origin.auto_apply_eligible() || rec.destructive {
1133 return false;
1134 }
1135 let manifest_ok = self
1139 .analyzers
1140 .iter()
1141 .map(|a| a.manifest())
1142 .find(|m| m.id == rec.analyzer)
1143 .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
1144 if !manifest_ok {
1145 return false;
1146 }
1147 let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1148 return false;
1149 };
1150 let family = crate::manifest::analyzer_family(&rec.analyzer);
1151 if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1152 return false;
1153 }
1154 match &rec.proposal {
1159 Proposal::Cal { cal } => cal
1160 .lines()
1161 .map(str::trim)
1162 .filter(|l| !l.is_empty())
1163 .all(|l| supersede_is_value_identical(sub, l)),
1164 _ => false,
1165 }
1166 }
1167
1168 fn auto_apply<S: OmsSubstrate>(
1171 &self,
1172 sub: &mut S,
1173 p: &mut LoopPersisted,
1174 rec: &Recommendation,
1175 now_ms: i64,
1176 ) -> Result<()> {
1177 let mut created = Vec::new();
1178 if let Proposal::Cal { cal } = &rec.proposal {
1179 if cal.lines().map(str::trim).any(is_definition_statement) {
1184 return Err(Error::InvalidProposal(
1185 "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1186 auto-applied: it changes what every future context contains, so it \
1187 requires a human APPROVE + APPLY with BECAUSE"
1188 .into(),
1189 ));
1190 }
1191 for r in sub.execute_cal(cal)? {
1192 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1193 created.push(h.to_string());
1194 }
1195 }
1196 }
1197 let applied = AppliedRecord {
1198 applied_at_ms: now_ms,
1199 target_ref: rec.target_ref.clone(),
1200 rollbackable: rec.rollbackable,
1201 created_hashes: created,
1202 inverse_cal: None,
1203 metric: rec.metric.clone(),
1204 };
1205 let prev = p.audit_heads.get(&rec.hash).cloned();
1206 let audit = AuditRecord {
1207 rec_hash: rec.hash.clone(),
1208 from: Some(RecStatus::Pending),
1209 to: RecStatus::Applied,
1210 actor: "policy:auto".into(),
1211 observer_type: ObserverType::Policy,
1212 because: "auto-applied per host policy".into(),
1213 previous_audit_hash: prev,
1214 gating: None,
1215 at_ms: now_ms,
1216 };
1217 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1218 p.audit_heads.insert(rec.hash.clone(), audit_hash);
1219 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1220 p.applied.insert(rec.hash.clone(), applied);
1221 Ok(())
1222 }
1223
1224 #[allow(clippy::too_many_arguments)]
1227 pub fn review<S: OmsSubstrate>(
1228 &self,
1229 sub: &mut S,
1230 rec_hash: &str,
1231 decision: Decision,
1232 actor: &str,
1233 observer: ObserverType,
1234 scopes: &ScopeSet,
1235 because: &str,
1236 now_ms: i64,
1237 ) -> Result<()> {
1238 if !scopes.has(Scope::Review) {
1239 return Err(Error::ScopeDenied("review".into()));
1240 }
1241 let because = validate_because(because)?;
1242 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1243 let status = *p
1244 .status_index
1245 .get(rec_hash)
1246 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1247 let to = match decision {
1248 Decision::Approve => RecStatus::Approved,
1249 Decision::Reject => RecStatus::Rejected,
1250 };
1251 if !status.can_transition_to(to, false) {
1252 return Err(Error::LifecycleViolation(format!(
1253 "{} -> {}",
1254 status.as_str(),
1255 to.as_str()
1256 )));
1257 }
1258 if to == RecStatus::Approved {
1259 if let Some(creator) = p.creators.get(rec_hash) {
1260 if creator == actor {
1261 return Err(Error::SelfApproval(format!(
1262 "{actor} created this recommendation"
1263 )));
1264 }
1265 }
1266 if let Some(trigger) = p.co_creators.get(rec_hash) {
1267 if trigger == actor {
1268 return Err(Error::SelfApproval(format!(
1269 "{actor} triggered the run that authored this recommendation"
1270 )));
1271 }
1272 }
1273 }
1274 let prev = p.audit_heads.get(rec_hash).cloned();
1275 let audit = AuditRecord {
1276 rec_hash: rec_hash.into(),
1277 from: Some(status),
1278 to,
1279 actor: actor.into(),
1280 observer_type: observer,
1281 because,
1282 previous_audit_hash: prev,
1283 gating: None,
1284 at_ms: now_ms,
1285 };
1286 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1287 p.audit_heads.insert(rec_hash.into(), audit_hash);
1288 p.status_index.insert(rec_hash.into(), to);
1289 if to == RecStatus::Rejected {
1290 if let Ok(rec) = load_rec(sub, rec_hash) {
1291 strike_cooldown(&mut p, rec.dedup_key, now_ms);
1292 }
1293 }
1294 sub.store_state(&p.to_value()?)?;
1295 Ok(())
1296 }
1297
1298 pub fn preflight_apply<S: OmsSubstrate>(
1315 &self,
1316 sub: &S,
1317 rec_hash: &str,
1318 scopes: &ScopeSet,
1319 allow_destructive: bool,
1320 has_gating: bool,
1321 ) -> Result<()> {
1322 if !scopes.has(Scope::Apply) {
1323 return Err(Error::ScopeDenied("apply".into()));
1324 }
1325 let rec = load_rec(sub, rec_hash)?;
1326 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1327 return Err(Error::DestructiveGated(
1328 "destructive apply requires admin scope + allow_destructive".into(),
1329 ));
1330 }
1331 ensure_executable(rec.action_kind, &rec.proposal)?;
1332 if requires_gating(rec.action_kind) && !has_gating {
1333 return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1334 }
1335 Ok(())
1336 }
1337
1338 #[allow(clippy::too_many_arguments)]
1342 pub fn apply<S: OmsSubstrate>(
1343 &self,
1344 sub: &mut S,
1345 rec_hash: &str,
1346 actor: &str,
1347 observer: ObserverType,
1348 scopes: &ScopeSet,
1349 because: &str,
1350 allow_destructive: bool,
1351 now_ms: i64,
1352 ) -> Result<AppliedRecord> {
1353 self.apply_inner(
1354 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1355 )
1356 }
1357
1358 pub fn gating_evidence<S: OmsSubstrate>(
1365 &self,
1366 sub: &S,
1367 rec_hash: &str,
1368 run_id: &str,
1369 ) -> Result<crate::recommendation::GatingEvidence> {
1370 let rec = self
1371 .recommendations(sub, None)?
1372 .into_iter()
1373 .find(|r| r.hash == rec_hash)
1374 .ok_or_else(|| {
1375 Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1376 })?;
1377 let pin = rec.evalset_hash.ok_or_else(|| {
1378 Error::InvalidProposal(
1379 "this recommendation pins no evalset — a gating run applies only \
1380 to code and adapter revisions"
1381 .into(),
1382 )
1383 })?;
1384 match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1389 Some(run) => Ok(crate::recommendation::GatingEvidence {
1390 evalset_hash: pin,
1391 run_id: run.run_id,
1392 passed: run.passed,
1393 failed: run.failed,
1394 }),
1395 None => Err(Error::InvalidProposal(format!(
1396 "no recorded gate run '{run_id}' for evalset {pin} — run \
1397 `areev eval run --evalset {pin} ...` first"
1398 ))),
1399 }
1400 }
1401
1402 #[allow(clippy::too_many_arguments)]
1406 pub fn apply_gated<S: OmsSubstrate>(
1407 &self,
1408 sub: &mut S,
1409 rec_hash: &str,
1410 actor: &str,
1411 observer: ObserverType,
1412 scopes: &ScopeSet,
1413 because: &str,
1414 allow_destructive: bool,
1415 gating: &crate::recommendation::GatingEvidence,
1416 now_ms: i64,
1417 ) -> Result<AppliedRecord> {
1418 self.apply_inner(
1419 sub,
1420 rec_hash,
1421 actor,
1422 observer,
1423 scopes,
1424 because,
1425 allow_destructive,
1426 Some(gating),
1427 now_ms,
1428 )
1429 }
1430
1431 #[allow(clippy::too_many_arguments)]
1432 fn apply_inner<S: OmsSubstrate>(
1433 &self,
1434 sub: &mut S,
1435 rec_hash: &str,
1436 actor: &str,
1437 observer: ObserverType,
1438 scopes: &ScopeSet,
1439 because: &str,
1440 allow_destructive: bool,
1441 gating: Option<&crate::recommendation::GatingEvidence>,
1442 now_ms: i64,
1443 ) -> Result<AppliedRecord> {
1444 if !scopes.has(Scope::Apply) {
1445 return Err(Error::ScopeDenied("apply".into()));
1446 }
1447 let because = validate_because(because)?;
1448 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1449 let status = *p
1450 .status_index
1451 .get(rec_hash)
1452 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1453 if !status.can_transition_to(RecStatus::Applied, false) {
1454 return Err(Error::LifecycleViolation(format!(
1455 "{} -> applied (approve first)",
1456 status.as_str()
1457 )));
1458 }
1459 let rec = load_rec(sub, rec_hash)?;
1460 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1461 return Err(Error::DestructiveGated(
1462 "destructive apply requires admin scope + allow_destructive".into(),
1463 ));
1464 }
1465 if requires_gating(rec.action_kind) {
1470 let g = gating
1471 .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1472 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1473 if g.evalset_hash != pin {
1474 return Err(Error::InvalidProposal(format!(
1475 "gating ran evalset {} but the recommendation is pinned \
1476 to {pin} (Rule E1)",
1477 g.evalset_hash
1478 )));
1479 }
1480 match sub.grain(pin)? {
1481 Some(evalset) if evalset.is_live() => {}
1482 Some(_) => {
1483 return Err(Error::InvalidProposal(
1484 "the pinned evalset was superseded after gating — \
1485 the recommendation must re-gate (Rule E1)"
1486 .into(),
1487 ))
1488 }
1489 None => {
1490 return Err(Error::InvalidProposal(format!(
1491 "pinned evalset {pin} not found in the substrate"
1492 )))
1493 }
1494 }
1495 if g.failed > 0 {
1496 return Err(Error::InvalidProposal(format!(
1497 "the gating run failed {}/{} cases — a failing gate \
1498 admits nothing",
1499 g.failed,
1500 g.passed + g.failed
1501 )));
1502 }
1503 }
1504
1505 let mut created = Vec::new();
1507 let mut inverse_cal: Option<String> = None;
1511 match &rec.proposal {
1512 Proposal::Cal { cal } => {
1513 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1514 if !is_definition_statement(line) {
1515 continue;
1516 }
1517 match sub.definition_inverse(line)? {
1518 Some(inv) => inverse_cal = Some(inv),
1519 None => {
1520 return Err(Error::InvalidProposal(format!(
1521 "this substrate cannot record a rollback inverse for {line:?}; \
1522 a definition rewrite that ROLLBACK could not undo is refused \
1523 rather than applied"
1524 )))
1525 }
1526 }
1527 }
1528 let rows = sub.execute_cal(cal)?;
1529 for r in rows {
1530 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1531 created.push(h.to_string());
1532 }
1533 }
1534 }
1535 Proposal::Edit { .. } => {
1538 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1539 }
1540 Proposal::Data { data } if requires_gating(rec.action_kind) => {
1548 let relation = if rec.action_kind == ActionKind::AdapterRevision {
1549 "mg:adapter_promotion"
1550 } else {
1551 "mg:code_promotion"
1552 };
1553 let mut promoted = data.clone();
1561 if let Some(Value::String(src)) = promoted.remove("source") {
1562 let address = sub.put_blob(src.as_bytes())?;
1563 promoted.insert("code_address".into(), Value::from(address));
1564 }
1565 let mut spec = crate::substrate::GrainSpec::new(
1566 crate::model::grain_type::FACT,
1567 LOOP_NS,
1568 )
1569 .with_field("subject", rec.target_ref.clone())
1570 .with_field("relation", relation)
1571 .with_field(
1572 "object",
1573 serde_json::to_string(&promoted)
1574 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1575 )
1576 .with_field("rec_hash", rec_hash.to_string());
1577 if let Some(g) = gating {
1578 spec = spec
1579 .with_field("gating_evalset", g.evalset_hash.clone())
1580 .with_field("gating_run_id", g.run_id.clone());
1581 }
1582 created.push(sub.put_grain(&spec)?);
1583 }
1584 Proposal::Data { data } => {
1585 let revert_of = data
1590 .get("revert_of")
1591 .and_then(Value::as_str)
1592 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1593 self.rollback(
1594 sub,
1595 revert_of,
1596 actor,
1597 observer,
1598 scopes,
1599 &because,
1600 now_ms,
1601 )?;
1602 p = LoopPersisted::from_value(sub.load_state()?)?;
1605 if let Ok(reverted) = load_rec(sub, revert_of) {
1615 strike_cooldown(&mut p, reverted.dedup_key, now_ms);
1616 }
1617 }
1618 }
1619
1620 let applied = AppliedRecord {
1621 applied_at_ms: now_ms,
1622 target_ref: rec.target_ref.clone(),
1623 rollbackable: rec.rollbackable,
1624 created_hashes: created,
1625 inverse_cal,
1626 metric: rec.metric.clone(),
1627 };
1628 let prev = p.audit_heads.get(rec_hash).cloned();
1629 let audit = AuditRecord {
1630 rec_hash: rec_hash.into(),
1631 from: Some(status),
1632 to: RecStatus::Applied,
1633 actor: actor.into(),
1634 observer_type: observer,
1635 because,
1636 previous_audit_hash: prev,
1637 gating: gating.cloned(),
1638 at_ms: now_ms,
1639 };
1640 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1641 p.audit_heads.insert(rec_hash.into(), audit_hash);
1642 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1643 p.applied.insert(rec_hash.into(), applied.clone());
1644 sub.store_state(&p.to_value()?)?;
1645 Ok(applied)
1646 }
1647
1648 #[allow(clippy::too_many_arguments)]
1651 pub fn rollback<S: OmsSubstrate>(
1652 &self,
1653 sub: &mut S,
1654 rec_hash: &str,
1655 actor: &str,
1656 observer: ObserverType,
1657 scopes: &ScopeSet,
1658 because: &str,
1659 now_ms: i64,
1660 ) -> Result<()> {
1661 if !scopes.has(Scope::Apply) {
1662 return Err(Error::ScopeDenied("apply".into()));
1663 }
1664 let because = validate_because(because)?;
1665 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1666 let status = *p
1667 .status_index
1668 .get(rec_hash)
1669 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1670 if !status.can_transition_to(RecStatus::RolledBack, false) {
1671 return Err(Error::LifecycleViolation(format!(
1672 "{} -> rolled_back",
1673 status.as_str()
1674 )));
1675 }
1676 let applied = p
1677 .applied
1678 .get(rec_hash)
1679 .cloned()
1680 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1681 if !applied.rollbackable {
1682 return Err(Error::LifecycleViolation(
1683 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1684 ));
1685 }
1686 for h in &applied.created_hashes {
1687 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1688 }
1689 if let Some(inverse) = &applied.inverse_cal {
1696 sub.execute_cal(inverse)?;
1697 }
1698 let prev = p.audit_heads.get(rec_hash).cloned();
1699 let audit = AuditRecord {
1700 rec_hash: rec_hash.into(),
1701 from: Some(status),
1702 to: RecStatus::RolledBack,
1703 actor: actor.into(),
1704 observer_type: observer,
1705 because,
1706 previous_audit_hash: prev,
1707 gating: None,
1708 at_ms: now_ms,
1709 };
1710 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1711 p.audit_heads.insert(rec_hash.into(), audit_hash);
1712 p.status_index
1713 .insert(rec_hash.into(), RecStatus::RolledBack);
1714 sub.store_state(&p.to_value()?)?;
1715 Ok(())
1716 }
1717
1718 pub fn recommendations<S: OmsSubstrate>(
1723 &self,
1724 sub: &S,
1725 status_filter: Option<RecStatus>,
1726 ) -> Result<Vec<Recommendation>> {
1727 let p = LoopPersisted::from_value(sub.load_state()?)?;
1728 let grains = sub.grains_of_type(
1729 crate::model::grain_type::RECOMMENDATION,
1730 Some(LOOP_NS),
1731 ReadOpts {
1732 live_only: false,
1733 since_ms: None,
1734 },
1735 )?;
1736 let mut out = Vec::new();
1737 for g in grains {
1738 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1739 rec.status = p
1740 .status_index
1741 .get(&g.hash)
1742 .copied()
1743 .unwrap_or(RecStatus::Pending);
1744 if let Some(f) = status_filter {
1745 if rec.status != f {
1746 continue;
1747 }
1748 }
1749 out.push(rec);
1750 }
1751 out.sort_by(|a, b| {
1760 b.severity
1761 .cmp(&a.severity)
1762 .then(a.created_at_ms.cmp(&b.created_at_ms))
1763 .then(a.dedup_key.cmp(&b.dedup_key))
1764 .then(a.hash.cmp(&b.hash))
1765 });
1766 Ok(out)
1767 }
1768
1769 pub fn analyzer_settings<S: OmsSubstrate>(
1772 &self,
1773 sub: &S,
1774 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1775 let p = LoopPersisted::from_value(sub.load_state()?)?;
1776 Ok(self
1777 .analyzers
1778 .iter()
1779 .map(|a| {
1780 let m = a.manifest();
1781 let cfg = p.config.get(&m.id);
1782 crate::config::AnalyzerSetting {
1783 id: m.id.clone(),
1784 title: m.title.clone(),
1785 description: m.description.clone(),
1786 tier: format!("{:?}", m.tier),
1787 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1788 default_on: m.default_on,
1789 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1790 severity_floor: cfg
1791 .and_then(|c| c.severity_floor)
1792 .map(|s| s.as_str().to_string()),
1793 }
1794 })
1795 .collect())
1796 }
1797
1798 pub fn set_analyzer_config<S: OmsSubstrate>(
1804 &self,
1805 sub: &mut S,
1806 analyzer_id: &str,
1807 update: crate::config::AnalyzerConfigUpdate,
1808 scopes: &ScopeSet,
1809 ) -> Result<crate::config::AnalyzerConfig> {
1810 if !scopes.has(Scope::Admin) {
1811 return Err(Error::ScopeDenied("admin".into()));
1812 }
1813 let manifest = self
1814 .analyzers
1815 .iter()
1816 .map(|a| a.manifest())
1817 .find(|m| m.id == analyzer_id)
1818 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1819 if let Some(params) = &update.params {
1821 manifest.resolve_params(params)?;
1822 }
1823 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1824 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1825 if let Some(enabled) = update.enabled {
1826 cfg.enabled = Some(enabled);
1827 }
1828 if update.clear_floor {
1829 cfg.severity_floor = None;
1830 } else if let Some(floor) = update.severity_floor {
1831 cfg.severity_floor = Some(floor);
1832 }
1833 if let Some(params) = update.params {
1834 cfg.params = params;
1835 }
1836 if let Some(ns) = update.namespaces {
1837 cfg.namespaces = ns;
1838 }
1839 let stored = cfg.clone();
1840 sub.store_state(&p.to_value()?)?;
1841 Ok(stored)
1842 }
1843
1844 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1847 let p = LoopPersisted::from_value(sub.load_state()?)?;
1848 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1849 out.sort_by(|a, b| {
1853 a.measured_at_ms
1854 .cmp(&b.measured_at_ms)
1855 .then(a.horizon_ms.cmp(&b.horizon_ms))
1856 .then(a.metric.cmp(&b.metric))
1857 .then(a.rec_hash.cmp(&b.rec_hash))
1858 });
1859 Ok(out)
1860 }
1861
1862 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1866 let p = LoopPersisted::from_value(sub.load_state()?)?;
1867 let new = count_new(sub, p.state.watermark_ms)?;
1868 let (grains_since_run, error_events_since_run) = (new.grains, new.error_events);
1869 let recs = self.recommendations(sub, None)?;
1870 let mut pending = 0;
1871 let mut applied = 0;
1872 for r in &recs {
1873 match r.status {
1874 RecStatus::Pending => pending += 1,
1875 RecStatus::Applied => applied += 1,
1876 _ => {}
1877 }
1878 }
1879 let stale = match p.state.last_run_ms {
1881 None => true,
1882 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1883 };
1884 Ok(Health {
1885 last_run_ms: p.state.last_run_ms,
1886 grains_since_run,
1887 error_events_since_run,
1888 pending,
1889 applied,
1890 total: recs.len() as u64,
1891 stale,
1892 })
1893 }
1894
1895 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1900 let recs = self.recommendations(sub, None)?;
1901 let mut m = LlmMetrics::default();
1902 for r in &recs {
1903 if !matches!(r.origin, Origin::Llm { .. }) {
1904 continue;
1905 }
1906 m.proposed += 1;
1907 match r.status {
1908 RecStatus::Pending => m.pending += 1,
1909 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1910 RecStatus::Rejected => m.rejected += 1,
1911 RecStatus::Expired => {}
1912 }
1913 }
1914 let decided = m.approved + m.rejected;
1915 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1916 Ok(m)
1917 }
1918}
1919
1920#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1922pub struct Health {
1923 #[serde(skip_serializing_if = "Option::is_none")]
1924 pub last_run_ms: Option<i64>,
1925 pub grains_since_run: u64,
1926 pub error_events_since_run: u64,
1927 pub pending: u64,
1928 pub applied: u64,
1929 pub total: u64,
1930 pub stale: bool,
1933}
1934
1935#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1937pub struct LlmMetrics {
1938 pub proposed: u64,
1941 pub pending: u64,
1942 pub approved: u64,
1944 pub rejected: u64,
1945 #[serde(skip_serializing_if = "Option::is_none")]
1947 pub approval_rate: Option<f64>,
1948}
1949
1950fn measure_outcomes<S: OmsSubstrate>(
1957 sub: &S,
1958 p: &mut LoopPersisted,
1959 now_ms: i64,
1960) -> Result<Vec<OutcomeInput>> {
1961 let mut due: Vec<(String, crate::config::AppliedRecord, Checkpoint)> = Vec::new();
1965 for (h, a) in &p.applied {
1966 if p.status_index.get(h) != Some(&RecStatus::Applied) {
1967 continue;
1968 }
1969 let Some(metric) = &a.metric else { continue };
1970 let done = p.measured.get(h).cloned().unwrap_or_default();
1971 for cp in metric.schedule() {
1972 if !done.contains(&cp) && checkpoint_due(sub, metric, a.applied_at_ms, cp, now_ms)? {
1973 due.push((h.clone(), a.clone(), cp));
1974 }
1975 }
1976 }
1977
1978 let mut out = Vec::new();
1979 for (rec_hash, applied, checkpoint) in due {
1980 let metric = applied.metric.as_ref().unwrap();
1981 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1982 continue; };
1984 let baseline = baseline_at_apply(sub, metric, applied.applied_at_ms)?;
1985 let regressed = crate::recommendation::is_regression(
1986 baseline,
1987 current,
1988 metric.higher_is_better,
1989 );
1990 p.outcomes.entry(rec_hash.clone()).or_default().push(
1991 crate::recommendation::OutcomeResult {
1992 rec_hash: rec_hash.clone(),
1993 metric: metric.metric.clone(),
1994 baseline,
1995 current,
1996 verdict: if regressed { "regressed" } else { "held" }.into(),
1997 horizon_ms: checkpoint.as_ms().unwrap_or(0),
2000 checkpoint: checkpoint.as_ms().is_none().then_some(checkpoint),
2001 measured_at_ms: now_ms,
2002 },
2003 );
2004 p.measured.entry(rec_hash.clone()).or_default().push(checkpoint);
2005 if regressed {
2006 out.push(OutcomeInput {
2007 rec_hash,
2008 target_ref: applied.target_ref.clone(),
2009 metric: metric.metric.clone(),
2010 baseline,
2011 current,
2012 unit: metric.unit.clone(),
2013 higher_is_better: metric.higher_is_better,
2014 });
2015 }
2016 }
2017 Ok(out)
2018}
2019
2020fn checkpoint_due<S: SubstrateRead>(
2029 sub: &S,
2030 metric: &crate::recommendation::MetricSnapshot,
2031 applied_at_ms: i64,
2032 cp: Checkpoint,
2033 now_ms: i64,
2034) -> Result<bool> {
2035 Ok(match cp {
2036 Checkpoint::AfterMs(ms) => now_ms - applied_at_ms >= ms,
2037 Checkpoint::AfterRuns(n) => match crate::eval::parse_evalset_metric(&metric.metric) {
2038 Some((evalset, _)) => {
2039 crate::eval::eval_runs(sub, evalset, Some(applied_at_ms + 1))?.len() >= n as usize
2040 }
2041 None => false,
2042 },
2043 Checkpoint::AfterGrains(n) => count_new(sub, Some(applied_at_ms))?.grains >= n as u64,
2044 })
2045}
2046
2047fn baseline_at_apply<S: SubstrateRead>(
2059 sub: &S,
2060 metric: &crate::recommendation::MetricSnapshot,
2061 applied_at_ms: i64,
2062) -> Result<f64> {
2063 if let Some((evalset, field)) = crate::eval::parse_evalset_metric(&metric.metric) {
2064 if let Some(run) = crate::eval::newest_eval_run_before(sub, evalset, applied_at_ms)? {
2065 if let Some(v) = crate::eval::run_value(&run, field) {
2066 return Ok(v);
2067 }
2068 }
2069 }
2070 Ok(metric.baseline)
2071}
2072
2073pub(crate) fn measure_metric<S: SubstrateRead>(
2075 sub: &S,
2076 metric: &crate::recommendation::MetricSnapshot,
2077 since_ms: i64,
2078) -> Result<Option<f64>> {
2079 match metric.metric.as_str() {
2080 "tool_error_recurrence" => {
2085 let Some(tool) = &metric.subject else { return Ok(None) };
2086 let tools = sub.grains_of_type(
2087 crate::model::grain_type::TOOL,
2088 None,
2089 ReadOpts { live_only: true, since_ms: Some(since_ms) },
2090 )?;
2091 let n = tools
2092 .iter()
2093 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
2094 .filter(|t| {
2095 metric.relation.as_deref().is_none_or(|sig| {
2098 crate::analyzers::tool_failure::normalize_signature(
2099 t.tool_content().unwrap_or(""),
2100 ) == sig
2101 })
2102 })
2103 .count();
2104 Ok(Some(n as f64))
2105 }
2106 "contradiction_recurrence" => {
2110 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
2111 return Ok(None);
2112 };
2113 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
2114 let distinct: BTreeSet<String> = facts
2115 .iter()
2116 .filter(|f| {
2117 f.fact_relation()
2118 .is_some_and(|r| normalize_ident(r) == *relation)
2119 })
2120 .filter_map(|f| f.fact_object().map(normalize_ident))
2121 .collect();
2122 Ok(Some(distinct.len().saturating_sub(1) as f64))
2123 }
2124 m if m.starts_with("evalset:") => {
2137 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
2138 return Ok(None);
2139 };
2140 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
2141 return Ok(None);
2142 };
2143 Ok(crate::eval::run_value(&run, field))
2144 }
2145 _ => Ok(None),
2146 }
2147}
2148
2149fn scoped_live_facts<S: SubstrateRead>(
2152 sub: &S,
2153 namespace: Option<&str>,
2154 subject: &str,
2155) -> Result<Vec<GrainRecord>> {
2156 let facts = sub.grains_of_type(
2157 crate::model::grain_type::FACT,
2158 None,
2159 ReadOpts { live_only: true, since_ms: None },
2160 )?;
2161 Ok(facts
2162 .into_iter()
2163 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
2164 .filter(|f| {
2165 f.fact_subject()
2166 .is_some_and(|s| normalize_ident(s) == subject)
2167 })
2168 .collect())
2169}
2170
2171fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
2183 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
2184 return false;
2185 };
2186 if fields.is_empty() {
2187 return false;
2188 }
2189 let Ok(Some(grain)) = sub.grain(&target) else {
2190 return false;
2191 };
2192 if grain.valid_to_ms.is_some() {
2202 return false;
2203 }
2204 fields.iter().all(|(k, v)| {
2206 if k == "namespace" {
2207 return v
2208 .as_str()
2209 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2210 }
2211 match (v, grain.fields.get(k)) {
2212 (Value::String(a), Some(Value::String(b))) => {
2213 normalize_ident(a) == normalize_ident(b)
2214 }
2215 (a, Some(b)) => a == b,
2216 (_, None) => false,
2217 }
2218 })
2219}
2220
2221const EVIDENCE_CAP: usize = 64;
2231const CITED_SEED_CAP: usize = 24;
2232const TOOL_SEED_CAP: usize = 16;
2233const NOTE_SEED_CAP: usize = 8;
2238const LENS_RESERVE: usize = 24;
2239
2240const MIN_LLM_CONFIDENCE: f64 = 0.75;
2243
2244macro_rules! discover_instructions {
2254 ($scoring:literal) => {
2255 concat!(
2256 "You review an agent's memory for quality. \
2257Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2258findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2259stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2260recurring cost or hand-off the agent's own setup could remove). \
2261The deterministic findings already cover what the ERROR TEXT says; restating \
2262one of them earns nothing. The evidence may also contain OUTCOME records — a \
2263run's observable shape together with whether it was accepted or rejected. A \
2264problem that raised no error at all is exactly the kind the deterministic \
2265checks cannot see, so compare the rejected outcomes against the accepted \
2266ones: a feature they share and the accepted ones lack is a candidate rule. \
2267Require at least two rejected outcomes before proposing one — a single \
2268rejection is an anecdote, not a pattern. ",
2269 $scoring,
2270 " The 'approved' and 'rejected' lists, when \
2271present, show findings this reviewer recently accepted or rejected — prefer the \
2272kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2273or more evidence items from the bundle by their 'id' (or 'hash'), name a \
2274'target', and include your confidence 0.0-1.0. Return JSON: \
2275{\"recommendations\":[{\"summary\":\"...\",\
2276\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<id>\"],\
2277\"confidence\":0.0,\"proposal\":{...}}]}. \
2278OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2279you are not asking to change anything. Include it ONLY when the evidence \
2280supports a specific change, choosing exactly one kind: \
2281(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2282\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
2283an action the agent itself takes on the next occasion. Either ADD an action \
2284it is failing to take ('Record the vendor name and the amount on every \
2285invoice, not just the date') or ORDER one that goes wrong ('Refund a \
2286subscription before cancelling it; refunds on cancelled subscriptions are \
2287refused'). Name the action, not a check on it: 'validate', 'verify' and \
2288'ensure ... is correct' describe a review step the agent has no way to \
2289perform, and such a rule changes nothing even once applied. \
2290(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2291entity target — a durable fact the agent keeps having to be told (an alias, a \
2292settled default, a preference). 'relation' is a short identifier (letters, \
2293digits, _ - . :), not a sentence. \
2294(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2295\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2296assembles the agent's context, when the evidence shows it retrieves the wrong \
2297things. Give the FULL new body; it replaces the old one. \
2298(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2299with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2300workflow. Only these paths are editable: 'edges.<i>.cond', \
2301'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2302holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2303(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2304— full replacement source for that tool. It is applied only after a recorded \
2305evaluation run passes, so propose one only when the evidence shows the current \
2306code is the defect. \
2307The subject of a fact, the name of a query, the plan hash and the tool name \
2308all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2309becomes a change a human reviewer may apply, so it must be fully supported by \
2310the cited evidence. Propose nothing you cannot ground in the evidence."
2311 )
2312 };
2313}
2314
2315const DISCOVER_INSTRUCTIONS: &str = discover_instructions!(
2317 "SCORING: propose a finding ONLY if you \
2318are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2319useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2320earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2321when there is nothing worth flagging."
2322);
2323
2324const DISCOVER_LEARNER_INSTRUCTIONS: &str = discover_instructions!(
2326 "SCORING: you are the learning stage of a deployed agent, and what you \
2327propose now is what it will do differently next time — a lesson you withhold \
2328is a mistake it repeats. A correct, actionable proposal earns 1; a wrong or \
2329trivial one is penalized 1; returning nothing while the evidence holds a \
2330recurring failure, two or more rejected outcomes, or an instruction from a \
2331person is ALSO penalized 1. Abstain only when the evidence shows none of \
2332those. Prefer the one proposal that addresses the most frequent or most costly \
2333failure over several speculative ones, and report your confidence honestly — \
2334an independent verifier, not you, decides what survives."
2335);
2336
2337const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2343fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2344the facts it relies on are actually present in the cited evidence, NOT that its \
2345conclusion is stated verbatim. Decompose the finding into the factual claims it \
2346depends on. Mark supported=true when those facts are present in the evidence \
2347(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2348on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2349different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2350
2351const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2353each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2354never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2355SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2356'possible' findings with no concrete defect, and reject any claimed \
2357inconsistency or contradiction that is not backed by at least two actually \
2358conflicting facts in the cited evidence. (2) Context — does the finding \
2359correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2360finding when it names a genuine, specific problem grounded in its evidence and \
2361materially useful to a human reviewer; otherwise reject it, and default to \
2362keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2363redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2364grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2365KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2366{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2367
2368const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2370guidance note to help a human reviewer decide. Do not restate the finding. Return \
2371JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2372
2373fn push_evidence(
2379 evidence: &mut Vec<crate::llm::EvidenceItem>,
2380 bundle: &mut BTreeSet<String>,
2381 ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2382 g: &GrainRecord,
2383 attribution: crate::policy::EvidenceAttribution,
2384) {
2385 if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2386 ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2387 evidence.push(crate::llm::EvidenceItem {
2388 id: format!("e{}", evidence.len() + 1),
2389 hash: g.hash.clone(),
2390 grain_type: g.grain_type.clone(),
2391 text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
2392 });
2393 }
2394}
2395
2396pub(crate) fn resolve_citation(
2404 cite: &str,
2405 bundle: &BTreeSet<String>,
2406 id_to_hash: &std::collections::BTreeMap<&str, &str>,
2407) -> Option<String> {
2408 let cite = cite.trim();
2409 if bundle.contains(cite) {
2410 return Some(cite.to_string());
2411 }
2412 if let Some(h) = id_to_hash.get(cite) {
2413 return Some((*h).to_string());
2414 }
2415 const MIN_PREFIX: usize = 12;
2416 if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
2417 let lower = cite.to_ascii_lowercase();
2418 let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
2419 if let (Some(h), None) = (it.next(), it.next()) {
2420 return Some(h.clone());
2421 }
2422 }
2423 None
2424}
2425
2426fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
2429 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2430 return format!("{s} {r} {o}");
2431 }
2432 if let Some(t) = g.tool_name() {
2437 let status = if g.is_error() { "error" } else { "ok" };
2438 let out = g.tool_content().unwrap_or("");
2439 let input = match g.fields.get("input") {
2445 Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
2446 Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
2447 _ => String::new(),
2448 };
2449 return format!("tool {t}{input} {status}: {out}");
2450 }
2451 for key in ["content", "body", "text", "summary", "object"] {
2459 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2460 if v.is_empty() {
2461 continue;
2462 }
2463 if attribution == crate::policy::EvidenceAttribution::Anonymous {
2474 return v.to_string();
2475 }
2476 if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
2477 if !who.is_empty() {
2478 let kind = g
2479 .fields
2480 .get("observer_type")
2481 .and_then(|v| v.as_str())
2482 .unwrap_or("");
2483 let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
2484 let mut prefix = if kind == "human" {
2485 format!("{who} (a person) said")
2486 } else {
2487 format!("{who} observed")
2488 };
2489 if !about.is_empty() {
2490 prefix.push_str(&format!(" of {about}"));
2491 }
2492 return format!("{prefix}: {v}");
2493 }
2494 }
2495 return v.to_string();
2496 }
2497 }
2498 String::new()
2499}
2500
2501fn sanitize_lesson(s: &str) -> String {
2506 sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2507}
2508
2509fn sanitize_line(s: &str, max: usize) -> String {
2515 let cleaned: String =
2516 s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2517 crate::llm::cap(cleaned.trim(), max)
2518}
2519
2520fn sanitize_relation(s: &str) -> Option<String> {
2524 let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2525 if r.is_empty()
2526 || !r
2527 .chars()
2528 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2529 {
2530 return None;
2531 }
2532 Some(r)
2533}
2534
2535fn safe_definition_body(body: &str) -> bool {
2551 if body.contains('{') || body.contains('}') {
2552 return false;
2553 }
2554 !body
2555 .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2556 .any(|tok| {
2557 ["FORGET", "PURGE", "DROP", "DEFINE"]
2558 .iter()
2559 .any(|kw| tok.eq_ignore_ascii_case(kw))
2560 })
2561}
2562
2563fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2568 let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2569 match resolved {
2570 Some(r) => format!("{summary} {}", r.rendered),
2571 None => summary,
2572 }
2573}
2574
2575struct ValidatedDraft {
2581 draft: crate::llm::LlmDraft,
2582 target_ref: String,
2583 cited: Vec<String>,
2584 resolved: Option<ResolvedProposal>,
2585}
2586
2587struct ResolvedProposal {
2591 action: ActionKind,
2592 proposal: Proposal,
2593 rendered: String,
2596 summary_key: &'static str,
2597 summary_args: serde_json::Map<String, Value>,
2598 rollbackable: bool,
2599 evalset_hash: Option<String>,
2600 importance: f64,
2601 fact_fields: Option<serde_json::Map<String, Value>>,
2605}
2606
2607fn plan_edit_allowed(path: &str) -> bool {
2617 let seg: Vec<&str> = path.split('.').collect();
2618 match seg.as_slice() {
2619 ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2620 ["retries", node] => !node.is_empty(),
2621 _ => false,
2622 }
2623}
2624
2625fn plan_get(body: &Value, path: &str) -> Value {
2628 let mut cur = body;
2629 for seg in path.split('.') {
2630 cur = match cur {
2631 Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2632 Some(v) => v,
2633 None => return Value::Null,
2634 },
2635 Value::Object(o) => match o.get(seg) {
2636 Some(v) => v,
2637 None => return Value::Null,
2638 },
2639 _ => return Value::Null,
2640 };
2641 }
2642 cur.clone()
2643}
2644
2645fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2648 let segs: Vec<&str> = path.split('.').collect();
2649 let Some((last, parents)) = segs.split_last() else {
2650 return false;
2651 };
2652 let mut cur = body;
2653 for seg in parents {
2654 cur = match cur {
2655 Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2656 Some(v) => v,
2657 None => return false,
2658 },
2659 Value::Object(o) => match o.get_mut(*seg) {
2660 Some(v) => v,
2661 None => return false,
2662 },
2663 _ => return false,
2664 };
2665 }
2666 match cur {
2667 Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2668 Some(slot) => {
2669 *slot = to;
2670 true
2671 }
2672 None => false,
2673 },
2674 Value::Object(o) => {
2675 o.insert((*last).to_string(), to);
2676 true
2677 }
2678 _ => false,
2679 }
2680}
2681
2682fn plan_value_ok(path: &str, to: &Value) -> bool {
2687 let seg: Vec<&str> = path.split('.').collect();
2688 match seg.as_slice() {
2689 ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
2690 !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
2691 }),
2692 ["edges", _, "max_cycles"] | ["retries", _] => {
2693 to.as_u64().is_some_and(|n| n <= 1_000)
2694 }
2695 _ => false,
2696 }
2697}
2698
2699fn resolve_proposal<S: OmsSubstrate>(
2705 sub: &S,
2706 d: &crate::llm::LlmDraft,
2707 target: &TargetRef,
2708 cited: &[String],
2709 ns_by_hash: &std::collections::BTreeMap<String, String>,
2710 caps: Capabilities,
2711 policy: &crate::policy::Policy,
2712) -> Option<ResolvedProposal> {
2713 use crate::llm::DraftProposal as P;
2714 let (skills, plans) = (&policy.skills, &policy.plans);
2715 let mut args = serde_json::Map::new();
2716 match d.parsed_proposal()? {
2717 P::Plan { description, when_to_use, nodes, edges } => {
2719 if !plans.enabled || !caps.plans {
2720 return None;
2721 }
2722 let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
2723 derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
2724 args.insert("name".into(), Value::from(name.clone()));
2725 args.insert("nodes".into(), Value::from(n_nodes as u64));
2726 args.insert("edges".into(), Value::from(n_edges as u64));
2727 let skill_stmt = match &existing_skill {
2728 Some(h) => cal::supersede(h, "skill", &skill),
2729 None => cal::add("skill", &skill),
2730 };
2731 let plan_stmt = match &existing_plan {
2732 Some(h) => cal::supersede(h, "workflow", &workflow),
2733 None => cal::add("workflow", &workflow),
2734 };
2735 let (action, verb) = if existing_skill.is_some() || existing_plan.is_some() {
2736 (ActionKind::Revise, "revise")
2737 } else {
2738 (ActionKind::Record, "record")
2739 };
2740 Some(ResolvedProposal {
2741 action,
2742 proposal: Proposal::Cal { cal: cal::batch(&[skill_stmt, plan_stmt]) },
2743 rendered: format!(
2744 "Proposed plan to {verb}: \"{name}\" — {n_nodes} steps, {n_edges} edges; when: {}",
2745 skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
2746 ),
2747 summary_key: "llm.plan",
2748 summary_args: args,
2749 rollbackable: true,
2750 evalset_hash: None,
2751 importance: 0.65,
2752 fact_fields: None,
2753 })
2754 }
2755 P::Skill { description, when_to_use, steps } => {
2757 if !skills.enabled {
2758 return None;
2759 }
2760 let SkillFields { fields, name, n_steps, existing } =
2761 derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
2762 args.insert("name".into(), Value::from(name.clone()));
2763 args.insert("steps".into(), Value::from(n_steps as u64));
2764 let (action, cal, verb) = match existing {
2767 Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
2768 None => (ActionKind::Record, cal::add("skill", &fields), "record"),
2769 };
2770 Some(ResolvedProposal {
2771 action,
2772 proposal: Proposal::Cal { cal },
2773 rendered: format!(
2774 "Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
2775 fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
2776 ),
2777 summary_key: "llm.skill",
2778 summary_args: args,
2779 rollbackable: true,
2780 evalset_hash: None,
2781 importance: 0.6,
2782 fact_fields: None,
2783 })
2784 }
2785 P::Lesson { lesson } => {
2787 let lesson = sanitize_lesson(&lesson);
2788 if lesson.is_empty() {
2789 return None;
2790 }
2791 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
2792 args.insert("lesson".into(), Value::from(lesson.clone()));
2793 Some(ResolvedProposal {
2794 action: ActionKind::ClusterFailure,
2798 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2799 rendered: format!("Proposed lesson to record: \"{lesson}\""),
2800 summary_key: "llm.lesson",
2801 summary_args: args,
2802 rollbackable: true,
2803 evalset_hash: None,
2804 importance: 0.5,
2805 fact_fields: Some(fields),
2806 })
2807 }
2808 P::Fact { relation, object } => {
2810 let relation = sanitize_relation(&relation)?;
2811 let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
2812 if object.is_empty() {
2813 return None;
2814 }
2815 let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
2816 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
2817 args.insert("relation".into(), Value::from(relation.clone()));
2818 args.insert("object".into(), Value::from(object.clone()));
2819 Some(ResolvedProposal {
2820 action: ActionKind::Record,
2821 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2822 rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
2823 summary_key: "llm.fact",
2824 summary_args: args,
2825 rollbackable: true,
2826 evalset_hash: None,
2827 importance: 0.5,
2828 fact_fields: Some(fields),
2829 })
2830 }
2831 P::QueryRevision { body } => {
2833 let name = target.opaque();
2834 if name.is_empty()
2838 || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
2839 {
2840 return None;
2841 }
2842 let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
2843 if body.is_empty() || !safe_definition_body(&body) {
2844 return None;
2845 }
2846 let stmt = match target.scheme() {
2847 "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
2848 "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
2849 _ => return None,
2850 };
2851 sub.validate_cal(&stmt).ok()?;
2856 sub.definition_inverse(&stmt).ok().flatten()?;
2857 args.insert("name".into(), Value::from(name));
2858 args.insert("body".into(), Value::from(body.clone()));
2859 Some(ResolvedProposal {
2860 action: ActionKind::Revise,
2861 proposal: Proposal::Cal { cal: stmt },
2862 rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
2863 summary_key: "llm.query_revision",
2864 summary_args: args,
2865 rollbackable: true,
2866 evalset_hash: None,
2867 importance: 0.6,
2868 fact_fields: None,
2869 })
2870 }
2871 P::PlanRevision { edits } => {
2873 if !caps.plans
2874 || target.scheme() != "grain"
2875 || edits.is_empty()
2876 || edits.len() > crate::llm::MAX_PLAN_EDITS
2877 {
2878 return None;
2879 }
2880 let hash = target.opaque();
2881 let g = sub.grain(hash).ok().flatten()?;
2882 if g.grain_type != "workflow" || !g.is_live() {
2883 return None;
2884 }
2885 let mut body = Value::Object(g.fields.clone());
2886 let mut deltas = Vec::new();
2887 let nodes: std::collections::BTreeSet<String> = body
2888 .get("nodes")
2889 .and_then(Value::as_array)
2890 .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
2891 .unwrap_or_default();
2892 for e in &edits {
2893 if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
2894 return None;
2895 }
2896 if let Some(node) = e.path.strip_prefix("retries.") {
2901 if !nodes.contains(node) {
2902 return None;
2903 }
2904 }
2905 if plan_get(&body, &e.path) != e.from {
2908 return None;
2909 }
2910 if e.from == e.to {
2914 return None;
2915 }
2916 if !plan_set(&mut body, &e.path, e.to.clone()) {
2917 return None;
2918 }
2919 deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
2920 }
2921 sub.validate_plan(&body).ok()?;
2925 let Value::Object(fields) = body else {
2926 return None;
2927 };
2928 let stmt = cal::supersede(hash, "workflow", &fields);
2929 sub.validate_cal(&stmt).ok()?;
2934 args.insert("plan".into(), Value::from(hash));
2935 args.insert("edits".into(), Value::from(deltas.join("; ")));
2936 Some(ResolvedProposal {
2937 action: ActionKind::Revise,
2938 proposal: Proposal::Cal { cal: stmt },
2939 rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
2940 summary_key: "llm.plan_revision",
2941 summary_args: args,
2942 rollbackable: true,
2943 evalset_hash: None,
2944 importance: 0.7,
2945 fact_fields: None,
2946 })
2947 }
2948 P::CodeRevision { source } => {
2950 if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
2951 return None;
2952 }
2953 if source.chars().count() > crate::llm::MAX_CODE_LEN {
2954 return None;
2955 }
2956 let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
2960 let mut data = serde_json::Map::new();
2961 data.insert("tool".into(), Value::from(target.opaque()));
2962 data.insert("source".into(), Value::from(source.clone()));
2963 args.insert("tool".into(), Value::from(target.opaque()));
2964 args.insert("bytes".into(), Value::from(source.len() as u64));
2965 Some(ResolvedProposal {
2966 action: ActionKind::CodeRevision,
2967 proposal: Proposal::Data { data },
2968 rendered: format!(
2969 "Proposed new source for tool {} ({} bytes), gated by evalset {}",
2970 target.opaque(),
2971 source.len(),
2972 evalset
2973 ),
2974 summary_key: "llm.code_revision",
2975 summary_args: args,
2976 rollbackable: true,
2977 evalset_hash: Some(evalset),
2978 importance: 0.8,
2979 fact_fields: None,
2980 })
2981 }
2982 }
2983}
2984
2985fn stamp_llm(
2995 model: &str,
2996 d: &crate::llm::LlmDraft,
2997 target_ref: String,
2998 cited: Vec<String>,
2999 resolved: Option<ResolvedProposal>,
3000 confidence: f64,
3001 now_ms: i64,
3002) -> Recommendation {
3003 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3004 let guidance = if d.guidance.trim().is_empty() {
3005 None
3006 } else {
3007 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
3008 };
3009 let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
3010 Some(mut r) => {
3011 let content = match &r.fact_fields {
3015 Some(fields) => format!(
3016 "{} {}",
3017 fields.get("relation").and_then(Value::as_str).unwrap_or(""),
3018 fields.get("object").and_then(Value::as_str).unwrap_or("")
3019 ),
3020 None => match &r.proposal {
3021 Proposal::Cal { cal } => cal.clone(),
3022 Proposal::Data { data } => Value::Object(data.clone()).to_string(),
3023 Proposal::Edit { diff, .. } => diff.clone(),
3024 },
3025 };
3026 if let Some(mut fields) = r.fact_fields.take() {
3029 fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
3030 r.proposal = Proposal::Cal { cal: cal::add("fact", &fields) };
3031 }
3032 let mut args = r.summary_args;
3033 args.insert("text".into(), Value::from(summary_text));
3034 (
3035 r.action,
3036 r.proposal,
3037 Summary::new(r.summary_key, args),
3038 r.rollbackable,
3039 r.importance,
3040 r.evalset_hash,
3041 Some(content),
3042 )
3043 }
3044 None => {
3045 let mut args = serde_json::Map::new();
3046 args.insert("text".into(), Value::from(summary_text));
3047 let mut data = serde_json::Map::new();
3048 data.insert("source".into(), Value::from("llm"));
3049 (
3050 ActionKind::Flag,
3051 Proposal::Data { data },
3052 Summary::new("llm.discover", args),
3053 false,
3054 0.3,
3055 None,
3056 None,
3057 )
3058 }
3059 };
3060 let dedup = match &content {
3064 Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
3065 None => dedup_key("llm", &target_ref, action),
3066 };
3067 Recommendation {
3068 hash: String::new(),
3069 analyzer: "loop.llm/1".to_string(),
3070 params_snapshot: serde_json::Map::new(),
3071 origin: Origin::Llm { model: model.to_string() },
3072 target_ref: target_ref.clone(),
3073 action_kind: action,
3074 dedup_key: dedup,
3075 summary,
3076 severity: Severity::Low,
3077 proposal,
3078 destructive: false,
3079 rollbackable,
3080 evidence: cited,
3081 evidence_query: None,
3082 metric: None,
3083 confidence: confidence.clamp(0.0, 1.0),
3085 importance,
3086 created_at_ms: now_ms,
3087 guidance,
3088 evalset_hash,
3089 status: RecStatus::Pending,
3090 }
3091}
3092
3093fn skill_instructions(min_steps: u32) -> String {
3103 format!(
3104 " (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
3105\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
3106PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
3107calls that reached its goal, which a later session facing the same situation \
3108should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
3109the tool called and the values that mattered (the field checked, the tag set, the \
3110exact format produced), a one-line description, and 'when_to_use' — the situation \
3111that should trigger it. The skill-name is a short identifier (letters, digits, \
3112_ -). If a saved skill already covers this procedure, use ITS name so it is \
3113patched rather than duplicated. Do not propose a skill for a procedure that \
3114failed, or for one already saved and unchanged.",
3115 crate::llm::MAX_SKILL_STEPS
3116 )
3117}
3118
3119fn plan_instructions(min_nodes: u32) -> String {
3123 format!(
3124 " (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
3125\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
3126\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
3127with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
3128but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
3129'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
3130the cited evidence — and what the step does with it; and 'edges' from step to \
3131step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
3132'path != literal', 'path exists' or '!path', where path is dotted names and the \
3133literal is a JSON string, number, true, false or null — no other operators; state \
3134a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
3135earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
3136has branches or a loop; prefer a skill when it is a straight list. If a saved \
3137plan already covers this procedure, use ITS name so it is patched.",
3138 crate::llm::MAX_PLAN_NODES
3139 )
3140}
3141
3142struct PlanFields {
3146 skill: serde_json::Map<String, Value>,
3147 workflow: serde_json::Map<String, Value>,
3148 name: String,
3149 n_nodes: usize,
3150 n_edges: usize,
3151 existing_skill: Option<String>,
3152 existing_plan: Option<String>,
3153}
3154
3155#[allow(clippy::too_many_arguments)]
3165fn derived_plan_fields<S: SubstrateRead>(
3166 sub: &S,
3167 target: &TargetRef,
3168 description: &str,
3169 when_to_use: &str,
3170 nodes: &[crate::llm::PlanNodeDraft],
3171 edges: &[crate::llm::PlanEdgeDraft],
3172 cited: &[String],
3173 ns_by_hash: &std::collections::BTreeMap<String, String>,
3174 plans: &crate::policy::PlanAuthoring,
3175) -> Option<PlanFields> {
3176 if target.scheme() != "entity" {
3177 return None;
3178 }
3179 let name = sanitize_skill_name(
3180 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3181 )?;
3182 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3183 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3184 if description.is_empty() || when_to_use.is_empty() {
3185 return None;
3186 }
3187 if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
3188 return None;
3189 }
3190 let known_tools: BTreeSet<String> = cited
3192 .iter()
3193 .filter_map(|h| sub.grain(h).ok().flatten())
3194 .filter_map(|g| g.tool_name().map(normalize_ident))
3195 .collect();
3196 let mut ids: Vec<String> = Vec::new();
3197 let mut steps: Vec<String> = Vec::new();
3198 let mut seen: BTreeSet<String> = BTreeSet::new();
3199 for n in nodes {
3200 let id = sanitize_skill_name(&n.id)?;
3201 if !seen.insert(id.clone()) {
3202 return None; }
3204 let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
3205 if tool.is_empty() || !known_tools.contains(&normalize_ident(&tool)) {
3206 return None; }
3208 let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
3209 if step.is_empty() {
3210 return None;
3211 }
3212 steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
3213 ids.push(id);
3214 }
3215 let mut edge_vals: Vec<Value> = Vec::new();
3216 let mut flow_lines: Vec<String> = Vec::new();
3217 for e in edges {
3218 let src = sanitize_skill_name(&e.src)?;
3219 let dst = sanitize_skill_name(&e.dst)?;
3220 if !seen.contains(&src) || !seen.contains(&dst) {
3221 return None;
3222 }
3223 let mut ev = serde_json::Map::new();
3224 ev.insert("src".into(), Value::from(src.clone()));
3225 ev.insert("dst".into(), Value::from(dst.clone()));
3226 let mut label = format!("{src} → {dst}");
3227 if let Some(c) = e
3228 .cond
3229 .as_deref()
3230 .map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
3231 .filter(|c| !c.is_empty())
3232 {
3233 label.push_str(&format!(" if {c}"));
3234 ev.insert("cond".into(), Value::from(c));
3235 }
3236 if let Some(m) = e.max_cycles {
3237 if m == 0 || m > 100 {
3238 return None;
3239 }
3240 label.push_str(&format!(" (at most {m} times)"));
3241 ev.insert("max_cycles".into(), Value::from(m));
3242 }
3243 flow_lines.push(label);
3244 edge_vals.push(Value::Object(ev));
3245 }
3246 if edge_vals.len() > 4 * ids.len() {
3247 return None;
3248 }
3249 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3251 for h in cited {
3252 if let Some(ns) = ns_by_hash.get(h) {
3253 if !ns.is_empty() {
3254 *ns_counts.entry(ns.as_str()).or_default() += 1;
3255 }
3256 }
3257 }
3258 let ns = ns_counts
3259 .iter()
3260 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3261 .map(|(ns, _)| ns.to_string());
3262
3263 let mut workflow = serde_json::Map::new();
3266 workflow.insert("nodes".into(), Value::from(ids.clone()));
3267 workflow.insert("edges".into(), Value::Array(edge_vals));
3268 workflow.insert("name".into(), Value::from(name.clone()));
3269 if let Some(ns) = &ns {
3270 workflow.insert("namespace".into(), Value::from(ns.clone()));
3271 }
3272 sub.validate_plan(&Value::Object(workflow.clone())).ok()?;
3273
3274 let mut instructions = steps.join("\n");
3277 if !flow_lines.is_empty() {
3278 instructions.push_str("\n\nFlow:\n");
3279 instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
3280 }
3281 let mut skill = serde_json::Map::new();
3282 skill.insert("name".into(), Value::from(name.clone()));
3283 skill.insert("description".into(), Value::from(description));
3284 skill.insert("when_to_use".into(), Value::from(when_to_use));
3285 skill.insert("instructions".into(), Value::from(instructions));
3286 if let Some(ns) = &ns {
3287 skill.insert("namespace".into(), Value::from(ns.clone()));
3288 }
3289 let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
3290 sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
3291 .ok()?
3292 .into_iter()
3293 .find(|g| pick(g))
3294 .map(|g| g.hash)
3295 };
3296 let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
3297 let existing_plan = live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str()));
3298 Some(PlanFields {
3299 skill,
3300 workflow,
3301 name,
3302 n_nodes: ids.len(),
3303 n_edges: flow_lines.len(),
3304 existing_skill,
3305 existing_plan,
3306 })
3307}
3308
3309pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
3311
3312fn detect_premise_drift<S: OmsSubstrate>(
3328 sub: &S,
3329 p: &mut LoopPersisted,
3330 now_ms: i64,
3331) -> Result<Vec<OutcomeInput>> {
3332 let mut out = Vec::new();
3333 let applied: Vec<(String, String, Vec<String>)> = p
3334 .applied
3335 .iter()
3336 .filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
3337 .map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
3338 .collect();
3339 for (rec_hash, target_ref, own) in applied {
3340 let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
3341 if rec.evidence.is_empty() {
3342 continue;
3343 }
3344 let mut moved = 0u64;
3345 for e in &rec.evidence {
3346 match sub.grain(e)? {
3347 None => moved += 1, Some(g) => {
3349 let Some(newer) = &g.superseded_by else { continue };
3350 if own.iter().any(|c| c == newer) {
3351 continue; }
3353 match sub.grain(newer)? {
3354 None => moved += 1,
3357 Some(n) => {
3358 if !same_value(&g, &n) {
3359 moved += 1;
3360 }
3361 }
3362 }
3363 }
3364 }
3365 }
3366 if moved == 0 {
3367 continue;
3368 }
3369 let already = p
3370 .outcomes
3371 .get(&rec_hash)
3372 .and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
3373 .is_some_and(|o| o.current == moved as f64);
3374 if !already {
3375 p.outcomes.entry(rec_hash.clone()).or_default().push(
3376 crate::recommendation::OutcomeResult {
3377 rec_hash: rec_hash.clone(),
3378 metric: PREMISE_DRIFT_METRIC.into(),
3379 baseline: 0.0,
3380 current: moved as f64,
3381 verdict: "drifted".into(),
3382 horizon_ms: 0,
3383 checkpoint: None,
3384 measured_at_ms: now_ms,
3385 },
3386 );
3387 }
3388 out.push(OutcomeInput {
3389 rec_hash,
3390 target_ref,
3391 metric: PREMISE_DRIFT_METRIC.into(),
3392 baseline: 0.0,
3393 current: moved as f64,
3394 unit: "superseded premises".into(),
3395 higher_is_better: false,
3396 });
3397 }
3398 Ok(out)
3399}
3400
3401fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
3406 if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
3407 return normalize_ident(a) == normalize_ident(b);
3408 }
3409 for key in ["content", "tool_content", "body", "text", "object"] {
3410 if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
3411 return normalize_ident(a) == normalize_ident(b);
3412 }
3413 }
3414 false
3415}
3416
3417fn sanitize_skill_name(s: &str) -> Option<String> {
3419 let t = s.trim();
3420 if t.is_empty()
3421 || t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
3422 || !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
3423 {
3424 return None;
3425 }
3426 Some(t.to_string())
3427}
3428
3429struct SkillFields {
3433 fields: serde_json::Map<String, Value>,
3434 name: String,
3435 n_steps: usize,
3436 existing: Option<String>,
3437}
3438
3439#[allow(clippy::too_many_arguments)]
3444fn derived_skill_fields<S: SubstrateRead>(
3445 sub: &S,
3446 target: &TargetRef,
3447 description: &str,
3448 when_to_use: &str,
3449 steps: &[String],
3450 cited: &[String],
3451 ns_by_hash: &std::collections::BTreeMap<String, String>,
3452 skills: &crate::policy::SkillAuthoring,
3453) -> Option<SkillFields> {
3454 if target.scheme() != "entity" {
3455 return None;
3456 }
3457 let name = sanitize_skill_name(
3458 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3459 )?;
3460 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3461 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3462 let steps: Vec<String> = steps
3463 .iter()
3464 .map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
3465 .filter(|st| !st.is_empty())
3466 .take(crate::llm::MAX_SKILL_STEPS)
3467 .collect();
3468 if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
3469 return None;
3470 }
3471 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3474 for h in cited {
3475 if let Some(ns) = ns_by_hash.get(h) {
3476 if !ns.is_empty() {
3477 *ns_counts.entry(ns.as_str()).or_default() += 1;
3478 }
3479 }
3480 }
3481 let ns = ns_counts
3482 .iter()
3483 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3484 .map(|(ns, _)| ns.to_string());
3485 let instructions = steps
3486 .iter()
3487 .enumerate()
3488 .map(|(i, st)| format!("{}. {st}", i + 1))
3489 .collect::<Vec<_>>()
3490 .join("\n");
3491 let mut fields = serde_json::Map::new();
3492 fields.insert("name".into(), Value::from(name.clone()));
3493 fields.insert("description".into(), Value::from(description));
3494 fields.insert("when_to_use".into(), Value::from(when_to_use));
3495 fields.insert("instructions".into(), Value::from(instructions));
3496 if let Some(ns) = &ns {
3497 fields.insert("namespace".into(), Value::from(ns.clone()));
3498 }
3499 let existing = sub
3501 .grains_of_type(
3502 crate::model::grain_type::SKILL,
3503 ns.as_deref(),
3504 ReadOpts { live_only: true, since_ms: None },
3505 )
3506 .ok()?
3507 .into_iter()
3508 .find(|g| g.skill_name() == Some(name.as_str()))
3509 .map(|g| g.hash);
3510 Some(SkillFields { fields, name, n_steps: steps.len(), existing })
3511}
3512
3513fn derived_fact_fields(
3514 target: &TargetRef,
3515 relation: &str,
3516 object: &str,
3517 cited: &[String],
3518 ns_by_hash: &std::collections::BTreeMap<String, String>,
3519) -> Option<serde_json::Map<String, Value>> {
3520 if target.scheme() != "entity" {
3521 return None;
3522 }
3523 let subject = target
3524 .opaque()
3525 .rsplit_once('/')
3526 .map(|(_, s)| s)
3527 .unwrap_or(target.opaque());
3528 if subject.is_empty() {
3529 return None;
3530 }
3531 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3532 for h in cited {
3533 if let Some(ns) = ns_by_hash.get(h) {
3534 if !ns.is_empty() {
3535 *ns_counts.entry(ns.as_str()).or_default() += 1;
3536 }
3537 }
3538 }
3539 let lesson_ns = ns_counts
3540 .iter()
3541 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3542 .map(|(ns, _)| ns.to_string());
3543 let mut fields = serde_json::Map::new();
3544 fields.insert("subject".into(), Value::from(subject));
3545 fields.insert("relation".into(), Value::from(relation));
3546 fields.insert("object".into(), Value::from(object));
3547 if let Some(ns) = lesson_ns {
3551 fields.insert("namespace".into(), Value::from(ns));
3552 }
3553 Some(fields)
3554}
3555
3556fn requires_gating(kind: ActionKind) -> bool {
3561 matches!(
3562 kind,
3563 ActionKind::CodeRevision | ActionKind::AdapterRevision
3564 )
3565}
3566
3567fn stamp(
3568 m: &AnalyzerManifest,
3569 params: &crate::manifest::Params,
3570 d: crate::recommendation::RecDraft,
3571 now_ms: i64,
3572) -> Result<Recommendation> {
3573 let target = TargetRef::parse(&d.target_ref)?;
3574 crate::recommendation::validate_code_rules(
3578 d.action_kind,
3579 target.target_class(),
3580 d.evalset_hash.as_deref(),
3581 )?;
3582 let revert_of = match (&d.action_kind, &d.proposal) {
3585 (ActionKind::Revert, Proposal::Data { data }) => {
3586 data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
3587 }
3588 _ => None,
3589 };
3590 let dedup = match revert_of.as_deref() {
3591 Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
3592 None => dedup_key(m.family(), &d.target_ref, d.action_kind),
3593 };
3594 let destructive = match &d.proposal {
3595 Proposal::Cal { cal } => cal::contains_destructive(cal),
3596 _ => false,
3597 };
3598 let rollbackable = match &d.proposal {
3599 Proposal::Cal { .. } => !destructive,
3600 Proposal::Edit { .. } => false,
3601 Proposal::Data { .. } => requires_gating(d.action_kind),
3605 };
3606 let mut evidence = d.evidence;
3607 evidence.truncate(MAX_EVIDENCE);
3608 let origin = match m.trust_class {
3613 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
3614 _ => Origin::Builtin,
3615 };
3616 Ok(Recommendation {
3617 hash: String::new(),
3618 analyzer: m.id.clone(),
3619 params_snapshot: params.snapshot(),
3620 origin,
3621 target_ref: target.as_string(),
3622 action_kind: d.action_kind,
3623 dedup_key: dedup,
3624 summary: d.summary,
3625 severity: d.severity,
3626 proposal: d.proposal,
3627 destructive,
3628 rollbackable,
3629 evidence,
3630 evidence_query: d.evidence_query,
3631 metric: d.metric,
3632 confidence: d.confidence,
3633 importance: d.importance,
3634 created_at_ms: now_ms,
3635 guidance: None,
3636 evalset_hash: d.evalset_hash,
3637 status: RecStatus::Pending,
3638 })
3639}
3640
3641fn validate_because(because: &str) -> Result<String> {
3642 let trimmed = because.trim();
3643 if trimmed.is_empty() {
3644 return Err(Error::InvalidProposal(
3645 "a BECAUSE reason is required".into(),
3646 ));
3647 }
3648 if trimmed.chars().count() > MAX_BECAUSE {
3649 return Err(Error::InvalidProposal(format!(
3650 "BECAUSE exceeds {MAX_BECAUSE} chars"
3651 )));
3652 }
3653 Ok(trimmed.to_string())
3654}
3655
3656fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
3657 for req in &m.requires {
3658 match req {
3659 Capability::Forks if !caps.forks => return Some("forks"),
3660 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
3661 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
3662 _ => {}
3663 }
3664 }
3665 None
3666}
3667
3668fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
3669 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
3670}
3671
3672fn gate(
3673 opts: &RunOptions,
3674 p: &LoopPersisted,
3675 new_grains: u64,
3676 new_errors: u64,
3677 now_ms: i64,
3678) -> Option<SkipReason> {
3679 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
3680 if !any {
3681 return None;
3682 }
3683 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
3684 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
3685 let stale_ok = opts
3686 .if_stale_ms
3687 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
3688 if min_new_ok || min_err_ok || stale_ok {
3689 return None;
3690 }
3691 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
3693 Some(SkipReason::NotStale)
3694 } else {
3695 Some(SkipReason::MinNewNotMet)
3696 }
3697}
3698
3699#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
3701pub(crate) struct NewSince {
3702 pub grains: u64,
3704 pub error_events: u64,
3706 pub events: u64,
3708 pub sessions: u64,
3710}
3711
3712fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
3713 let opts = ReadOpts {
3714 live_only: false,
3715 since_ms: watermark.map(|w| w + 1),
3716 };
3717 let mut n = NewSince::default();
3718 let mut sessions: BTreeSet<&str> = BTreeSet::new();
3719 let mut events_held: Vec<GrainRecord> = Vec::new();
3720 for t in [
3721 crate::model::grain_type::FACT,
3722 crate::model::grain_type::EVENT,
3723 crate::model::grain_type::TOOL,
3724 crate::model::grain_type::OBSERVATION,
3725 ] {
3726 let g = sub.grains_of_type(t, None, opts)?;
3727 n.grains += g.len() as u64;
3728 if t == crate::model::grain_type::TOOL {
3730 n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
3731 }
3732 if t == crate::model::grain_type::EVENT {
3733 n.events = g.len() as u64;
3734 events_held = g;
3735 }
3736 }
3737 for e in &events_held {
3738 if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
3739 sessions.insert(sid);
3740 }
3741 }
3742 n.sessions = sessions.len() as u64;
3743 Ok(n)
3744}
3745
3746fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
3753 if !c.is_set() {
3754 return None;
3755 }
3756 let time_ok = c
3757 .every_ms
3758 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
3759 let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
3760 let events_ok = c.every_events.is_some_and(|m| new.events >= m);
3761 let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
3762 if time_ok || grains_ok || events_ok || sessions_ok {
3763 None
3764 } else {
3765 Some(SkipReason::CadenceNotDue)
3766 }
3767}
3768
3769fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
3770 let grains = sub.grains_of_type(
3771 crate::model::grain_type::RECOMMENDATION,
3772 Some(LOOP_NS),
3773 ReadOpts {
3774 live_only: false,
3775 since_ms: None,
3776 },
3777 )?;
3778 let mut set = BTreeSet::new();
3779 for g in grains {
3780 let status = p
3781 .status_index
3782 .get(&g.hash)
3783 .copied()
3784 .unwrap_or(RecStatus::Pending);
3785 if matches!(
3792 status,
3793 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
3794 ) {
3795 if let Some(key) = g.str_field("dedup_key") {
3796 set.insert(key.to_string());
3797 }
3798 }
3799 }
3800 Ok(set)
3801}
3802
3803fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
3810 const BASE_MS: i64 = 7 * 86_400_000;
3811 const CAP_MS: i64 = 90 * 86_400_000;
3812 let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
3813 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
3814 *strikes = strikes.saturating_add(1);
3815 p.cooldowns.insert(dedup_key, now_ms + interval);
3816}
3817
3818fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
3819 let g = sub
3820 .grain(rec_hash)?
3821 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
3822 Recommendation::from_fields(rec_hash, &g.fields)
3823}
3824
3825pub(crate) fn is_definition_statement(line: &str) -> bool {
3833 let up = line.trim_start().to_ascii_uppercase();
3834 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
3835}
3836
3837const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
3840 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
3841 approve it to acknowledge it and let it expire.";
3842
3843const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
3848 (evalset hash + run id + stats) — use apply_gated";
3849
3850const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
3851 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
3852 acknowledge it and let it expire.";
3853
3854pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
3867 match proposal {
3868 Proposal::Cal { .. } => Ok(()),
3869 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
3870 Proposal::Data { data } => {
3876 if requires_gating(action_kind)
3877 || data.get("revert_of").and_then(Value::as_str).is_some()
3878 {
3879 Ok(())
3880 } else {
3881 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
3882 }
3883 }
3884 }
3885}
3886
3887#[cfg(test)]
3888mod definition_body_tests {
3889 use super::safe_definition_body;
3890
3891 #[test]
3892 fn ordinary_bodies_pass() {
3893 assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
3894 assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
3895 }
3896
3897 #[test]
3898 fn a_body_cannot_close_its_own_block_or_carry_destruction() {
3899 assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
3904 assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
3905 assert!(!safe_definition_body("RECALL facts FORGET abc"));
3907 assert!(!safe_definition_body("recall facts purge older than 1d"));
3908 assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
3909 assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
3911 assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
3913 }
3914}
3915
3916#[cfg(test)]
3917mod plan_edit_tests {
3918 use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
3919 use serde_json::json;
3920
3921 fn plan() -> serde_json::Value {
3922 json!({
3923 "nodes": ["fetch", "review", "post"],
3924 "edges": [
3925 {"src": "fetch", "dst": "review"},
3926 {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
3927 ],
3928 "bindings": {"fetch": "sha256:tool1"},
3929 "retries": {"fetch": 1}
3930 })
3931 }
3932
3933 #[test]
3934 fn the_allowlist_admits_thresholds_and_refuses_topology() {
3935 assert!(plan_edit_allowed("edges.1.cond"));
3936 assert!(plan_edit_allowed("edges.1.max_cycles"));
3937 assert!(plan_edit_allowed("retries.fetch"));
3938 for path in [
3941 "nodes",
3942 "nodes.0",
3943 "edges.0.src",
3944 "edges.0.dst",
3945 "edges",
3946 "bindings.fetch",
3947 "edges.x.cond",
3948 "",
3949 ] {
3950 assert!(!plan_edit_allowed(path), "{path} must not be editable");
3951 }
3952 }
3953
3954 #[test]
3955 fn values_are_type_checked_against_the_field() {
3956 assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
3959 assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
3960 assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
3961 assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
3962 assert!(plan_value_ok("retries.fetch", &json!(3)));
3963 assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
3964 assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
3965 assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
3966 assert!(!plan_value_ok("edges.0.src", &json!("other")));
3967 }
3968
3969 #[test]
3970 fn get_reads_through_arrays_and_objects_and_absence_is_null() {
3971 let p = plan();
3972 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
3973 assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
3974 assert_eq!(plan_get(&p, "retries.review"), json!(null));
3977 assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
3978 assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
3979 }
3980
3981 #[test]
3982 fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
3983 let mut p = plan();
3984 assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
3985 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
3986 assert!(plan_set(&mut p, "retries.review", json!(2)));
3987 assert_eq!(plan_get(&p, "retries.review"), json!(2));
3988 assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
3989 assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
3990 }
3991}
3992
3993#[cfg(test)]
3994mod definition_proposal_tests {
3995 use super::is_definition_statement;
3996
3997 #[test]
3998 fn definition_statements_are_recognized_in_both_spellings() {
3999 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
4000 assert!(is_definition_statement(" define template foo AS { x }"));
4001 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
4002 assert!(!is_definition_statement("ADD fact {}"));
4004 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
4005 assert!(!is_definition_statement("FORGET abc"));
4006 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
4010 }
4011}