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, an instruction from a \
2331person, or a multi-step procedure the agent completed successfully that no \
2332saved skill or plan covers, is ALSO penalized 1. Abstain only when the \
2333evidence shows none of those. Prefer the one proposal that addresses the most \
2334frequent or most costly failure — or, when nothing failed, the procedure that \
2335worked — over several speculative ones, and report your confidence honestly — \
2336an independent verifier, not you, decides what survives."
2337);
2338
2339const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2345fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2346the facts it relies on are actually present in the cited evidence, NOT that its \
2347conclusion is stated verbatim. Decompose the finding into the factual claims it \
2348depends on. Mark supported=true when those facts are present in the evidence \
2349(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2350on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2351different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2352
2353const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2355each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2356never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2357SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2358'possible' findings with no concrete defect, and reject any claimed \
2359inconsistency or contradiction that is not backed by at least two actually \
2360conflicting facts in the cited evidence. (2) Context — does the finding \
2361correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2362finding when it names a genuine, specific problem grounded in its evidence and \
2363materially useful to a human reviewer; otherwise reject it, and default to \
2364keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2365redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2366grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2367KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2368{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2369
2370const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2372guidance note to help a human reviewer decide. Do not restate the finding. Return \
2373JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2374
2375fn push_evidence(
2381 evidence: &mut Vec<crate::llm::EvidenceItem>,
2382 bundle: &mut BTreeSet<String>,
2383 ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2384 g: &GrainRecord,
2385 attribution: crate::policy::EvidenceAttribution,
2386) {
2387 if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2388 ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2389 evidence.push(crate::llm::EvidenceItem {
2390 id: format!("e{}", evidence.len() + 1),
2391 hash: g.hash.clone(),
2392 grain_type: g.grain_type.clone(),
2393 text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
2394 });
2395 }
2396}
2397
2398pub(crate) fn resolve_citation(
2406 cite: &str,
2407 bundle: &BTreeSet<String>,
2408 id_to_hash: &std::collections::BTreeMap<&str, &str>,
2409) -> Option<String> {
2410 let cite = cite.trim();
2411 if bundle.contains(cite) {
2412 return Some(cite.to_string());
2413 }
2414 if let Some(h) = id_to_hash.get(cite) {
2415 return Some((*h).to_string());
2416 }
2417 const MIN_PREFIX: usize = 12;
2418 if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
2419 let lower = cite.to_ascii_lowercase();
2420 let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
2421 if let (Some(h), None) = (it.next(), it.next()) {
2422 return Some(h.clone());
2423 }
2424 }
2425 None
2426}
2427
2428fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
2431 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2432 return format!("{s} {r} {o}");
2433 }
2434 if let Some(t) = g.tool_name() {
2439 let status = if g.is_error() { "error" } else { "ok" };
2440 let out = g.tool_content().unwrap_or("");
2441 let input = match g.fields.get("input") {
2447 Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
2448 Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
2449 _ => String::new(),
2450 };
2451 return format!("tool {t}{input} {status}: {out}");
2452 }
2453 for key in ["content", "body", "text", "summary", "object"] {
2461 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2462 if v.is_empty() {
2463 continue;
2464 }
2465 if attribution == crate::policy::EvidenceAttribution::Anonymous {
2476 return v.to_string();
2477 }
2478 if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
2479 if !who.is_empty() {
2480 let kind = g
2481 .fields
2482 .get("observer_type")
2483 .and_then(|v| v.as_str())
2484 .unwrap_or("");
2485 let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
2486 let mut prefix = if kind == "human" {
2487 format!("{who} (a person) said")
2488 } else {
2489 format!("{who} observed")
2490 };
2491 if !about.is_empty() {
2492 prefix.push_str(&format!(" of {about}"));
2493 }
2494 return format!("{prefix}: {v}");
2495 }
2496 }
2497 return v.to_string();
2498 }
2499 }
2500 String::new()
2501}
2502
2503fn sanitize_lesson(s: &str) -> String {
2508 sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2509}
2510
2511fn sanitize_line(s: &str, max: usize) -> String {
2517 let cleaned: String =
2518 s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2519 crate::llm::cap(cleaned.trim(), max)
2520}
2521
2522fn sanitize_relation(s: &str) -> Option<String> {
2526 let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2527 if r.is_empty()
2528 || !r
2529 .chars()
2530 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2531 {
2532 return None;
2533 }
2534 Some(r)
2535}
2536
2537fn safe_definition_body(body: &str) -> bool {
2553 if body.contains('{') || body.contains('}') {
2554 return false;
2555 }
2556 !body
2557 .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2558 .any(|tok| {
2559 ["FORGET", "PURGE", "DROP", "DEFINE"]
2560 .iter()
2561 .any(|kw| tok.eq_ignore_ascii_case(kw))
2562 })
2563}
2564
2565fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2570 let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2571 match resolved {
2572 Some(r) => format!("{summary} {}", r.rendered),
2573 None => summary,
2574 }
2575}
2576
2577struct ValidatedDraft {
2583 draft: crate::llm::LlmDraft,
2584 target_ref: String,
2585 cited: Vec<String>,
2586 resolved: Option<ResolvedProposal>,
2587}
2588
2589struct ResolvedProposal {
2593 action: ActionKind,
2594 proposal: Proposal,
2595 rendered: String,
2598 summary_key: &'static str,
2599 summary_args: serde_json::Map<String, Value>,
2600 rollbackable: bool,
2601 evalset_hash: Option<String>,
2602 importance: f64,
2603 fact_fields: Option<serde_json::Map<String, Value>>,
2607}
2608
2609fn plan_edit_allowed(path: &str) -> bool {
2619 let seg: Vec<&str> = path.split('.').collect();
2620 match seg.as_slice() {
2621 ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2622 ["retries", node] => !node.is_empty(),
2623 _ => false,
2624 }
2625}
2626
2627fn plan_get(body: &Value, path: &str) -> Value {
2630 let mut cur = body;
2631 for seg in path.split('.') {
2632 cur = match cur {
2633 Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2634 Some(v) => v,
2635 None => return Value::Null,
2636 },
2637 Value::Object(o) => match o.get(seg) {
2638 Some(v) => v,
2639 None => return Value::Null,
2640 },
2641 _ => return Value::Null,
2642 };
2643 }
2644 cur.clone()
2645}
2646
2647fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2650 let segs: Vec<&str> = path.split('.').collect();
2651 let Some((last, parents)) = segs.split_last() else {
2652 return false;
2653 };
2654 let mut cur = body;
2655 for seg in parents {
2656 cur = match cur {
2657 Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2658 Some(v) => v,
2659 None => return false,
2660 },
2661 Value::Object(o) => match o.get_mut(*seg) {
2662 Some(v) => v,
2663 None => return false,
2664 },
2665 _ => return false,
2666 };
2667 }
2668 match cur {
2669 Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2670 Some(slot) => {
2671 *slot = to;
2672 true
2673 }
2674 None => false,
2675 },
2676 Value::Object(o) => {
2677 o.insert((*last).to_string(), to);
2678 true
2679 }
2680 _ => false,
2681 }
2682}
2683
2684fn plan_value_ok(path: &str, to: &Value) -> bool {
2689 let seg: Vec<&str> = path.split('.').collect();
2690 match seg.as_slice() {
2691 ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
2692 !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
2693 }),
2694 ["edges", _, "max_cycles"] | ["retries", _] => {
2695 to.as_u64().is_some_and(|n| n <= 1_000)
2696 }
2697 _ => false,
2698 }
2699}
2700
2701fn resolve_proposal<S: OmsSubstrate>(
2707 sub: &S,
2708 d: &crate::llm::LlmDraft,
2709 target: &TargetRef,
2710 cited: &[String],
2711 ns_by_hash: &std::collections::BTreeMap<String, String>,
2712 caps: Capabilities,
2713 policy: &crate::policy::Policy,
2714) -> Option<ResolvedProposal> {
2715 use crate::llm::DraftProposal as P;
2716 let (skills, plans) = (&policy.skills, &policy.plans);
2717 let mut args = serde_json::Map::new();
2718 match d.parsed_proposal()? {
2719 P::Plan { description, when_to_use, nodes, edges } => {
2721 if !plans.enabled || !caps.plans {
2722 return None;
2723 }
2724 let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
2725 derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
2726 args.insert("name".into(), Value::from(name.clone()));
2727 args.insert("nodes".into(), Value::from(n_nodes as u64));
2728 let mut stmts = vec![match &existing_skill {
2729 Some(h) => cal::supersede(h, "skill", &skill),
2730 None => cal::add("skill", &skill),
2731 }];
2732 let (summary_key, kind) = match &workflow {
2739 Some(wf) => {
2740 stmts.push(match &existing_plan {
2741 Some(h) => cal::supersede(h, "workflow", wf),
2742 None => cal::add("workflow", wf),
2743 });
2744 args.insert("edges".into(), Value::from(n_edges as u64));
2745 ("llm.plan", "plan")
2746 }
2747 None => {
2748 args.insert("steps".into(), Value::from(n_nodes as u64));
2749 ("llm.skill", "skill")
2750 }
2751 };
2752 let patched = existing_skill.is_some() || (workflow.is_some() && existing_plan.is_some());
2753 let (action, verb) = if patched {
2754 (ActionKind::Revise, "revise")
2755 } else {
2756 (ActionKind::Record, "record")
2757 };
2758 Some(ResolvedProposal {
2759 action,
2760 proposal: Proposal::Cal { cal: cal::batch(&stmts) },
2761 rendered: format!(
2762 "Proposed {kind} to {verb}: \"{name}\" — {n_nodes} steps; when: {}",
2763 skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
2764 ),
2765 summary_key,
2766 summary_args: args,
2767 rollbackable: true,
2768 evalset_hash: None,
2769 importance: 0.65,
2770 fact_fields: None,
2771 })
2772 }
2773 P::Skill { description, when_to_use, steps } => {
2775 if !skills.enabled {
2776 return None;
2777 }
2778 let SkillFields { fields, name, n_steps, existing } =
2779 derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
2780 args.insert("name".into(), Value::from(name.clone()));
2781 args.insert("steps".into(), Value::from(n_steps as u64));
2782 let (action, cal, verb) = match existing {
2785 Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
2786 None => (ActionKind::Record, cal::add("skill", &fields), "record"),
2787 };
2788 Some(ResolvedProposal {
2789 action,
2790 proposal: Proposal::Cal { cal },
2791 rendered: format!(
2792 "Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
2793 fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
2794 ),
2795 summary_key: "llm.skill",
2796 summary_args: args,
2797 rollbackable: true,
2798 evalset_hash: None,
2799 importance: 0.6,
2800 fact_fields: None,
2801 })
2802 }
2803 P::Lesson { lesson } => {
2805 let lesson = sanitize_lesson(&lesson);
2806 if lesson.is_empty() {
2807 return None;
2808 }
2809 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
2810 args.insert("lesson".into(), Value::from(lesson.clone()));
2811 Some(ResolvedProposal {
2812 action: ActionKind::ClusterFailure,
2816 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2817 rendered: format!("Proposed lesson to record: \"{lesson}\""),
2818 summary_key: "llm.lesson",
2819 summary_args: args,
2820 rollbackable: true,
2821 evalset_hash: None,
2822 importance: 0.5,
2823 fact_fields: Some(fields),
2824 })
2825 }
2826 P::Fact { relation, object } => {
2828 let relation = sanitize_relation(&relation)?;
2829 let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
2830 if object.is_empty() {
2831 return None;
2832 }
2833 let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
2834 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
2835 args.insert("relation".into(), Value::from(relation.clone()));
2836 args.insert("object".into(), Value::from(object.clone()));
2837 Some(ResolvedProposal {
2838 action: ActionKind::Record,
2839 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2840 rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
2841 summary_key: "llm.fact",
2842 summary_args: args,
2843 rollbackable: true,
2844 evalset_hash: None,
2845 importance: 0.5,
2846 fact_fields: Some(fields),
2847 })
2848 }
2849 P::QueryRevision { body } => {
2851 let name = target.opaque();
2852 if name.is_empty()
2856 || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
2857 {
2858 return None;
2859 }
2860 let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
2861 if body.is_empty() || !safe_definition_body(&body) {
2862 return None;
2863 }
2864 let stmt = match target.scheme() {
2865 "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
2866 "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
2867 _ => return None,
2868 };
2869 sub.validate_cal(&stmt).ok()?;
2874 sub.definition_inverse(&stmt).ok().flatten()?;
2875 args.insert("name".into(), Value::from(name));
2876 args.insert("body".into(), Value::from(body.clone()));
2877 Some(ResolvedProposal {
2878 action: ActionKind::Revise,
2879 proposal: Proposal::Cal { cal: stmt },
2880 rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
2881 summary_key: "llm.query_revision",
2882 summary_args: args,
2883 rollbackable: true,
2884 evalset_hash: None,
2885 importance: 0.6,
2886 fact_fields: None,
2887 })
2888 }
2889 P::PlanRevision { edits } => {
2891 if !caps.plans
2892 || target.scheme() != "grain"
2893 || edits.is_empty()
2894 || edits.len() > crate::llm::MAX_PLAN_EDITS
2895 {
2896 return None;
2897 }
2898 let hash = target.opaque();
2899 let g = sub.grain(hash).ok().flatten()?;
2900 if g.grain_type != "workflow" || !g.is_live() {
2901 return None;
2902 }
2903 let mut body = Value::Object(g.fields.clone());
2904 let mut deltas = Vec::new();
2905 let nodes: std::collections::BTreeSet<String> = body
2906 .get("nodes")
2907 .and_then(Value::as_array)
2908 .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
2909 .unwrap_or_default();
2910 for e in &edits {
2911 if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
2912 return None;
2913 }
2914 if let Some(node) = e.path.strip_prefix("retries.") {
2919 if !nodes.contains(node) {
2920 return None;
2921 }
2922 }
2923 if plan_get(&body, &e.path) != e.from {
2926 return None;
2927 }
2928 if e.from == e.to {
2932 return None;
2933 }
2934 if !plan_set(&mut body, &e.path, e.to.clone()) {
2935 return None;
2936 }
2937 deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
2938 }
2939 sub.validate_plan(&body).ok()?;
2943 let Value::Object(fields) = body else {
2944 return None;
2945 };
2946 let stmt = cal::supersede(hash, "workflow", &fields);
2947 sub.validate_cal(&stmt).ok()?;
2952 args.insert("plan".into(), Value::from(hash));
2953 args.insert("edits".into(), Value::from(deltas.join("; ")));
2954 Some(ResolvedProposal {
2955 action: ActionKind::Revise,
2956 proposal: Proposal::Cal { cal: stmt },
2957 rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
2958 summary_key: "llm.plan_revision",
2959 summary_args: args,
2960 rollbackable: true,
2961 evalset_hash: None,
2962 importance: 0.7,
2963 fact_fields: None,
2964 })
2965 }
2966 P::CodeRevision { source } => {
2968 if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
2969 return None;
2970 }
2971 if source.chars().count() > crate::llm::MAX_CODE_LEN {
2972 return None;
2973 }
2974 let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
2978 let mut data = serde_json::Map::new();
2979 data.insert("tool".into(), Value::from(target.opaque()));
2980 data.insert("source".into(), Value::from(source.clone()));
2981 args.insert("tool".into(), Value::from(target.opaque()));
2982 args.insert("bytes".into(), Value::from(source.len() as u64));
2983 Some(ResolvedProposal {
2984 action: ActionKind::CodeRevision,
2985 proposal: Proposal::Data { data },
2986 rendered: format!(
2987 "Proposed new source for tool {} ({} bytes), gated by evalset {}",
2988 target.opaque(),
2989 source.len(),
2990 evalset
2991 ),
2992 summary_key: "llm.code_revision",
2993 summary_args: args,
2994 rollbackable: true,
2995 evalset_hash: Some(evalset),
2996 importance: 0.8,
2997 fact_fields: None,
2998 })
2999 }
3000 }
3001}
3002
3003fn stamp_llm(
3013 model: &str,
3014 d: &crate::llm::LlmDraft,
3015 target_ref: String,
3016 cited: Vec<String>,
3017 resolved: Option<ResolvedProposal>,
3018 confidence: f64,
3019 now_ms: i64,
3020) -> Recommendation {
3021 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3022 let guidance = if d.guidance.trim().is_empty() {
3023 None
3024 } else {
3025 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
3026 };
3027 let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
3028 Some(mut r) => {
3029 let content = match &r.fact_fields {
3033 Some(fields) => format!(
3034 "{} {}",
3035 fields.get("relation").and_then(Value::as_str).unwrap_or(""),
3036 fields.get("object").and_then(Value::as_str).unwrap_or("")
3037 ),
3038 None => match &r.proposal {
3039 Proposal::Cal { cal } => cal.clone(),
3040 Proposal::Data { data } => Value::Object(data.clone()).to_string(),
3041 Proposal::Edit { diff, .. } => diff.clone(),
3042 },
3043 };
3044 if let Some(mut fields) = r.fact_fields.take() {
3047 fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
3048 r.proposal = Proposal::Cal { cal: cal::add("fact", &fields) };
3049 }
3050 let mut args = r.summary_args;
3051 args.insert("text".into(), Value::from(summary_text));
3052 (
3053 r.action,
3054 r.proposal,
3055 Summary::new(r.summary_key, args),
3056 r.rollbackable,
3057 r.importance,
3058 r.evalset_hash,
3059 Some(content),
3060 )
3061 }
3062 None => {
3063 let mut args = serde_json::Map::new();
3064 args.insert("text".into(), Value::from(summary_text));
3065 let mut data = serde_json::Map::new();
3066 data.insert("source".into(), Value::from("llm"));
3067 (
3068 ActionKind::Flag,
3069 Proposal::Data { data },
3070 Summary::new("llm.discover", args),
3071 false,
3072 0.3,
3073 None,
3074 None,
3075 )
3076 }
3077 };
3078 let dedup = match &content {
3082 Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
3083 None => dedup_key("llm", &target_ref, action),
3084 };
3085 Recommendation {
3086 hash: String::new(),
3087 analyzer: "loop.llm/1".to_string(),
3088 params_snapshot: serde_json::Map::new(),
3089 origin: Origin::Llm { model: model.to_string() },
3090 target_ref: target_ref.clone(),
3091 action_kind: action,
3092 dedup_key: dedup,
3093 summary,
3094 severity: Severity::Low,
3095 proposal,
3096 destructive: false,
3097 rollbackable,
3098 evidence: cited,
3099 evidence_query: None,
3100 metric: None,
3101 confidence: confidence.clamp(0.0, 1.0),
3103 importance,
3104 created_at_ms: now_ms,
3105 guidance,
3106 evalset_hash,
3107 status: RecStatus::Pending,
3108 }
3109}
3110
3111fn skill_instructions(min_steps: u32) -> String {
3121 format!(
3122 " (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
3123\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
3124PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
3125calls that reached its goal, which a later session facing the same situation \
3126should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
3127the tool called and the values that mattered (the field checked, the tag set, the \
3128exact format produced), a one-line description, and 'when_to_use' — the situation \
3129that should trigger it. The skill-name is a short identifier (letters, digits, \
3130_ -). If a saved skill already covers this procedure, use ITS name so it is \
3131patched rather than duplicated. Do not propose a skill for a procedure that \
3132failed, or for one already saved and unchanged. A finding that itself describes \
3133two or more steps the agent should carry out in order ('after listing the \
3134tickets, fetch each, then …') IS a procedure: propose it as a skill or a plan, \
3135never as a lesson — a lesson is one rule, and a procedure written as one is a \
3136procedure nobody can open.",
3137 crate::llm::MAX_SKILL_STEPS
3138 )
3139}
3140
3141fn plan_instructions(min_nodes: u32) -> String {
3145 format!(
3146 " (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
3147\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
3148\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
3149with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
3150but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
3151'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
3152the cited evidence — and what the step does with it; and 'edges' from step to \
3153step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
3154'path != literal', 'path exists' or '!path', where path is dotted names and the \
3155literal is a JSON string, number, true, false or null — no other operators; state \
3156a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
3157earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
3158has branches or a loop; prefer a skill when it is a straight list. If a saved \
3159plan already covers this procedure, use ITS name so it is patched.",
3160 crate::llm::MAX_PLAN_NODES
3161 )
3162}
3163
3164struct PlanFields {
3168 skill: serde_json::Map<String, Value>,
3169 workflow: Option<serde_json::Map<String, Value>>,
3173 name: String,
3174 n_nodes: usize,
3175 n_edges: usize,
3176 existing_skill: Option<String>,
3177 existing_plan: Option<String>,
3178}
3179
3180#[allow(clippy::too_many_arguments)]
3190fn derived_plan_fields<S: SubstrateRead>(
3191 sub: &S,
3192 target: &TargetRef,
3193 description: &str,
3194 when_to_use: &str,
3195 nodes: &[crate::llm::PlanNodeDraft],
3196 edges: &[crate::llm::PlanEdgeDraft],
3197 cited: &[String],
3198 ns_by_hash: &std::collections::BTreeMap<String, String>,
3199 plans: &crate::policy::PlanAuthoring,
3200) -> Option<PlanFields> {
3201 if target.scheme() != "entity" {
3202 return None;
3203 }
3204 let name = sanitize_skill_name(
3205 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3206 )?;
3207 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3208 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3209 if description.is_empty() || when_to_use.is_empty() {
3210 return None;
3211 }
3212 if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
3213 return None;
3214 }
3215 let known_tools: BTreeSet<String> = cited
3217 .iter()
3218 .filter_map(|h| sub.grain(h).ok().flatten())
3219 .filter_map(|g| g.tool_name().map(normalize_ident))
3220 .collect();
3221 let mut ids: Vec<String> = Vec::new();
3222 let mut steps: Vec<String> = Vec::new();
3223 let mut seen: BTreeSet<String> = BTreeSet::new();
3224 let mut grounded = 0usize;
3225 for n in nodes {
3226 let id = sanitize_skill_name(&n.id)?;
3227 if !seen.insert(id.clone()) {
3228 return None; }
3230 let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
3231 if step.is_empty() {
3232 return None;
3233 }
3234 let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
3244 if !tool.is_empty() && known_tools.contains(&normalize_ident(&tool)) {
3245 grounded += 1;
3246 steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
3247 } else {
3248 steps.push(format!("{}. {id}: {step}", steps.len() + 1));
3249 }
3250 ids.push(id);
3251 }
3252 if grounded == 0 {
3255 return None;
3256 }
3257 let mut edge_vals: Vec<Value> = Vec::new();
3262 let mut flow_lines: Vec<String> = Vec::new();
3263 let mut runnable = true;
3264 for e in edges {
3265 let (Some(src), Some(dst)) = (sanitize_skill_name(&e.src), sanitize_skill_name(&e.dst)) else {
3266 runnable = false;
3267 continue;
3268 };
3269 if !seen.contains(&src) || !seen.contains(&dst) {
3270 flow_lines.push(format!("{} → {}", e.src.trim(), e.dst.trim()));
3272 runnable = false;
3273 continue;
3274 }
3275 let mut ev = serde_json::Map::new();
3276 ev.insert("src".into(), Value::from(src.clone()));
3277 ev.insert("dst".into(), Value::from(dst.clone()));
3278 let mut label = format!("{src} → {dst}");
3279 if let Some(c) = e
3280 .cond
3281 .as_deref()
3282 .map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
3283 .filter(|c| !c.is_empty())
3284 {
3285 label.push_str(&format!(" if {c}"));
3286 ev.insert("cond".into(), Value::from(c));
3287 }
3288 if let Some(m) = e.max_cycles {
3289 if m == 0 || m > 100 {
3290 runnable = false;
3291 } else {
3292 label.push_str(&format!(" (at most {m} times)"));
3293 ev.insert("max_cycles".into(), Value::from(m));
3294 }
3295 }
3296 flow_lines.push(label);
3297 edge_vals.push(Value::Object(ev));
3298 }
3299 if edge_vals.len() > 4 * ids.len() {
3300 runnable = false;
3301 }
3302 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3304 for h in cited {
3305 if let Some(ns) = ns_by_hash.get(h) {
3306 if !ns.is_empty() {
3307 *ns_counts.entry(ns.as_str()).or_default() += 1;
3308 }
3309 }
3310 }
3311 let ns = ns_counts
3312 .iter()
3313 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3314 .map(|(ns, _)| ns.to_string());
3315
3316 let mut workflow = serde_json::Map::new();
3319 workflow.insert("nodes".into(), Value::from(ids.clone()));
3320 workflow.insert("edges".into(), Value::Array(edge_vals));
3321 workflow.insert("name".into(), Value::from(name.clone()));
3322 if let Some(ns) = &ns {
3323 workflow.insert("namespace".into(), Value::from(ns.clone()));
3324 }
3325 let workflow = (runnable && sub.validate_plan(&Value::Object(workflow.clone())).is_ok())
3329 .then_some(workflow);
3330
3331 let mut instructions = steps.join("\n");
3334 if !flow_lines.is_empty() {
3335 instructions.push_str("\n\nFlow:\n");
3336 instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
3337 }
3338 let mut skill = serde_json::Map::new();
3339 skill.insert("name".into(), Value::from(name.clone()));
3340 skill.insert("description".into(), Value::from(description));
3341 skill.insert("when_to_use".into(), Value::from(when_to_use));
3342 skill.insert("instructions".into(), Value::from(instructions));
3343 if let Some(ns) = &ns {
3344 skill.insert("namespace".into(), Value::from(ns.clone()));
3345 }
3346 let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
3347 sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
3348 .ok()?
3349 .into_iter()
3350 .find(|g| pick(g))
3351 .map(|g| g.hash)
3352 };
3353 let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
3354 let existing_plan = workflow
3355 .is_some()
3356 .then(|| live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str())))
3357 .flatten();
3358 let n_edges = workflow
3359 .as_ref()
3360 .and_then(|w| w.get("edges"))
3361 .and_then(Value::as_array)
3362 .map_or(0, |a| a.len());
3363 Some(PlanFields { skill, workflow, name, n_nodes: ids.len(), n_edges, existing_skill, existing_plan })
3364}
3365
3366pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
3368
3369fn detect_premise_drift<S: OmsSubstrate>(
3385 sub: &S,
3386 p: &mut LoopPersisted,
3387 now_ms: i64,
3388) -> Result<Vec<OutcomeInput>> {
3389 let mut out = Vec::new();
3390 let applied: Vec<(String, String, Vec<String>)> = p
3391 .applied
3392 .iter()
3393 .filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
3394 .map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
3395 .collect();
3396 for (rec_hash, target_ref, own) in applied {
3397 let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
3398 if rec.evidence.is_empty() {
3399 continue;
3400 }
3401 let mut moved = 0u64;
3402 for e in &rec.evidence {
3403 match sub.grain(e)? {
3404 None => moved += 1, Some(g) => {
3406 let Some(newer) = &g.superseded_by else { continue };
3407 if own.iter().any(|c| c == newer) {
3408 continue; }
3410 match sub.grain(newer)? {
3411 None => moved += 1,
3414 Some(n) => {
3415 if !same_value(&g, &n) {
3416 moved += 1;
3417 }
3418 }
3419 }
3420 }
3421 }
3422 }
3423 if moved == 0 {
3424 continue;
3425 }
3426 let already = p
3427 .outcomes
3428 .get(&rec_hash)
3429 .and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
3430 .is_some_and(|o| o.current == moved as f64);
3431 if !already {
3432 p.outcomes.entry(rec_hash.clone()).or_default().push(
3433 crate::recommendation::OutcomeResult {
3434 rec_hash: rec_hash.clone(),
3435 metric: PREMISE_DRIFT_METRIC.into(),
3436 baseline: 0.0,
3437 current: moved as f64,
3438 verdict: "drifted".into(),
3439 horizon_ms: 0,
3440 checkpoint: None,
3441 measured_at_ms: now_ms,
3442 },
3443 );
3444 }
3445 out.push(OutcomeInput {
3446 rec_hash,
3447 target_ref,
3448 metric: PREMISE_DRIFT_METRIC.into(),
3449 baseline: 0.0,
3450 current: moved as f64,
3451 unit: "superseded premises".into(),
3452 higher_is_better: false,
3453 });
3454 }
3455 Ok(out)
3456}
3457
3458fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
3463 if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
3464 return normalize_ident(a) == normalize_ident(b);
3465 }
3466 for key in ["content", "tool_content", "body", "text", "object"] {
3467 if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
3468 return normalize_ident(a) == normalize_ident(b);
3469 }
3470 }
3471 false
3472}
3473
3474fn sanitize_skill_name(s: &str) -> Option<String> {
3476 let t = s.trim();
3477 if t.is_empty()
3478 || t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
3479 || !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
3480 {
3481 return None;
3482 }
3483 Some(t.to_string())
3484}
3485
3486struct SkillFields {
3490 fields: serde_json::Map<String, Value>,
3491 name: String,
3492 n_steps: usize,
3493 existing: Option<String>,
3494}
3495
3496#[allow(clippy::too_many_arguments)]
3501fn derived_skill_fields<S: SubstrateRead>(
3502 sub: &S,
3503 target: &TargetRef,
3504 description: &str,
3505 when_to_use: &str,
3506 steps: &[String],
3507 cited: &[String],
3508 ns_by_hash: &std::collections::BTreeMap<String, String>,
3509 skills: &crate::policy::SkillAuthoring,
3510) -> Option<SkillFields> {
3511 if target.scheme() != "entity" {
3512 return None;
3513 }
3514 let name = sanitize_skill_name(
3515 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3516 )?;
3517 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3518 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3519 let steps: Vec<String> = steps
3520 .iter()
3521 .map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
3522 .filter(|st| !st.is_empty())
3523 .take(crate::llm::MAX_SKILL_STEPS)
3524 .collect();
3525 if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
3526 return None;
3527 }
3528 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3531 for h in cited {
3532 if let Some(ns) = ns_by_hash.get(h) {
3533 if !ns.is_empty() {
3534 *ns_counts.entry(ns.as_str()).or_default() += 1;
3535 }
3536 }
3537 }
3538 let ns = ns_counts
3539 .iter()
3540 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3541 .map(|(ns, _)| ns.to_string());
3542 let instructions = steps
3543 .iter()
3544 .enumerate()
3545 .map(|(i, st)| format!("{}. {st}", i + 1))
3546 .collect::<Vec<_>>()
3547 .join("\n");
3548 let mut fields = serde_json::Map::new();
3549 fields.insert("name".into(), Value::from(name.clone()));
3550 fields.insert("description".into(), Value::from(description));
3551 fields.insert("when_to_use".into(), Value::from(when_to_use));
3552 fields.insert("instructions".into(), Value::from(instructions));
3553 if let Some(ns) = &ns {
3554 fields.insert("namespace".into(), Value::from(ns.clone()));
3555 }
3556 let existing = sub
3558 .grains_of_type(
3559 crate::model::grain_type::SKILL,
3560 ns.as_deref(),
3561 ReadOpts { live_only: true, since_ms: None },
3562 )
3563 .ok()?
3564 .into_iter()
3565 .find(|g| g.skill_name() == Some(name.as_str()))
3566 .map(|g| g.hash);
3567 Some(SkillFields { fields, name, n_steps: steps.len(), existing })
3568}
3569
3570fn derived_fact_fields(
3571 target: &TargetRef,
3572 relation: &str,
3573 object: &str,
3574 cited: &[String],
3575 ns_by_hash: &std::collections::BTreeMap<String, String>,
3576) -> Option<serde_json::Map<String, Value>> {
3577 if target.scheme() != "entity" {
3578 return None;
3579 }
3580 let subject = target
3581 .opaque()
3582 .rsplit_once('/')
3583 .map(|(_, s)| s)
3584 .unwrap_or(target.opaque());
3585 if subject.is_empty() {
3586 return None;
3587 }
3588 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3589 for h in cited {
3590 if let Some(ns) = ns_by_hash.get(h) {
3591 if !ns.is_empty() {
3592 *ns_counts.entry(ns.as_str()).or_default() += 1;
3593 }
3594 }
3595 }
3596 let lesson_ns = ns_counts
3597 .iter()
3598 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3599 .map(|(ns, _)| ns.to_string());
3600 let mut fields = serde_json::Map::new();
3601 fields.insert("subject".into(), Value::from(subject));
3602 fields.insert("relation".into(), Value::from(relation));
3603 fields.insert("object".into(), Value::from(object));
3604 if let Some(ns) = lesson_ns {
3608 fields.insert("namespace".into(), Value::from(ns));
3609 }
3610 Some(fields)
3611}
3612
3613fn requires_gating(kind: ActionKind) -> bool {
3618 matches!(
3619 kind,
3620 ActionKind::CodeRevision | ActionKind::AdapterRevision
3621 )
3622}
3623
3624fn stamp(
3625 m: &AnalyzerManifest,
3626 params: &crate::manifest::Params,
3627 d: crate::recommendation::RecDraft,
3628 now_ms: i64,
3629) -> Result<Recommendation> {
3630 let target = TargetRef::parse(&d.target_ref)?;
3631 crate::recommendation::validate_code_rules(
3635 d.action_kind,
3636 target.target_class(),
3637 d.evalset_hash.as_deref(),
3638 )?;
3639 let revert_of = match (&d.action_kind, &d.proposal) {
3642 (ActionKind::Revert, Proposal::Data { data }) => {
3643 data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
3644 }
3645 _ => None,
3646 };
3647 let dedup = match revert_of.as_deref() {
3648 Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
3649 None => dedup_key(m.family(), &d.target_ref, d.action_kind),
3650 };
3651 let destructive = match &d.proposal {
3652 Proposal::Cal { cal } => cal::contains_destructive(cal),
3653 _ => false,
3654 };
3655 let rollbackable = match &d.proposal {
3656 Proposal::Cal { .. } => !destructive,
3657 Proposal::Edit { .. } => false,
3658 Proposal::Data { .. } => requires_gating(d.action_kind),
3662 };
3663 let mut evidence = d.evidence;
3664 evidence.truncate(MAX_EVIDENCE);
3665 let origin = match m.trust_class {
3670 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
3671 _ => Origin::Builtin,
3672 };
3673 Ok(Recommendation {
3674 hash: String::new(),
3675 analyzer: m.id.clone(),
3676 params_snapshot: params.snapshot(),
3677 origin,
3678 target_ref: target.as_string(),
3679 action_kind: d.action_kind,
3680 dedup_key: dedup,
3681 summary: d.summary,
3682 severity: d.severity,
3683 proposal: d.proposal,
3684 destructive,
3685 rollbackable,
3686 evidence,
3687 evidence_query: d.evidence_query,
3688 metric: d.metric,
3689 confidence: d.confidence,
3690 importance: d.importance,
3691 created_at_ms: now_ms,
3692 guidance: None,
3693 evalset_hash: d.evalset_hash,
3694 status: RecStatus::Pending,
3695 })
3696}
3697
3698fn validate_because(because: &str) -> Result<String> {
3699 let trimmed = because.trim();
3700 if trimmed.is_empty() {
3701 return Err(Error::InvalidProposal(
3702 "a BECAUSE reason is required".into(),
3703 ));
3704 }
3705 if trimmed.chars().count() > MAX_BECAUSE {
3706 return Err(Error::InvalidProposal(format!(
3707 "BECAUSE exceeds {MAX_BECAUSE} chars"
3708 )));
3709 }
3710 Ok(trimmed.to_string())
3711}
3712
3713fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
3714 for req in &m.requires {
3715 match req {
3716 Capability::Forks if !caps.forks => return Some("forks"),
3717 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
3718 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
3719 _ => {}
3720 }
3721 }
3722 None
3723}
3724
3725fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
3726 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
3727}
3728
3729fn gate(
3730 opts: &RunOptions,
3731 p: &LoopPersisted,
3732 new_grains: u64,
3733 new_errors: u64,
3734 now_ms: i64,
3735) -> Option<SkipReason> {
3736 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
3737 if !any {
3738 return None;
3739 }
3740 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
3741 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
3742 let stale_ok = opts
3743 .if_stale_ms
3744 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
3745 if min_new_ok || min_err_ok || stale_ok {
3746 return None;
3747 }
3748 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
3750 Some(SkipReason::NotStale)
3751 } else {
3752 Some(SkipReason::MinNewNotMet)
3753 }
3754}
3755
3756#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
3758pub(crate) struct NewSince {
3759 pub grains: u64,
3761 pub error_events: u64,
3763 pub events: u64,
3765 pub sessions: u64,
3767}
3768
3769fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
3770 let opts = ReadOpts {
3771 live_only: false,
3772 since_ms: watermark.map(|w| w + 1),
3773 };
3774 let mut n = NewSince::default();
3775 let mut sessions: BTreeSet<&str> = BTreeSet::new();
3776 let mut events_held: Vec<GrainRecord> = Vec::new();
3777 for t in [
3778 crate::model::grain_type::FACT,
3779 crate::model::grain_type::EVENT,
3780 crate::model::grain_type::TOOL,
3781 crate::model::grain_type::OBSERVATION,
3782 ] {
3783 let g = sub.grains_of_type(t, None, opts)?;
3784 n.grains += g.len() as u64;
3785 if t == crate::model::grain_type::TOOL {
3787 n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
3788 }
3789 if t == crate::model::grain_type::EVENT {
3790 n.events = g.len() as u64;
3791 events_held = g;
3792 }
3793 }
3794 for e in &events_held {
3795 if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
3796 sessions.insert(sid);
3797 }
3798 }
3799 n.sessions = sessions.len() as u64;
3800 Ok(n)
3801}
3802
3803fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
3810 if !c.is_set() {
3811 return None;
3812 }
3813 let time_ok = c
3814 .every_ms
3815 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
3816 let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
3817 let events_ok = c.every_events.is_some_and(|m| new.events >= m);
3818 let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
3819 if time_ok || grains_ok || events_ok || sessions_ok {
3820 None
3821 } else {
3822 Some(SkipReason::CadenceNotDue)
3823 }
3824}
3825
3826fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
3827 let grains = sub.grains_of_type(
3828 crate::model::grain_type::RECOMMENDATION,
3829 Some(LOOP_NS),
3830 ReadOpts {
3831 live_only: false,
3832 since_ms: None,
3833 },
3834 )?;
3835 let mut set = BTreeSet::new();
3836 for g in grains {
3837 let status = p
3838 .status_index
3839 .get(&g.hash)
3840 .copied()
3841 .unwrap_or(RecStatus::Pending);
3842 if matches!(
3849 status,
3850 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
3851 ) {
3852 if let Some(key) = g.str_field("dedup_key") {
3853 set.insert(key.to_string());
3854 }
3855 }
3856 }
3857 Ok(set)
3858}
3859
3860fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
3867 const BASE_MS: i64 = 7 * 86_400_000;
3868 const CAP_MS: i64 = 90 * 86_400_000;
3869 let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
3870 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
3871 *strikes = strikes.saturating_add(1);
3872 p.cooldowns.insert(dedup_key, now_ms + interval);
3873}
3874
3875fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
3876 let g = sub
3877 .grain(rec_hash)?
3878 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
3879 Recommendation::from_fields(rec_hash, &g.fields)
3880}
3881
3882pub(crate) fn is_definition_statement(line: &str) -> bool {
3890 let up = line.trim_start().to_ascii_uppercase();
3891 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
3892}
3893
3894const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
3897 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
3898 approve it to acknowledge it and let it expire.";
3899
3900const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
3905 (evalset hash + run id + stats) — use apply_gated";
3906
3907const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
3908 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
3909 acknowledge it and let it expire.";
3910
3911pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
3924 match proposal {
3925 Proposal::Cal { .. } => Ok(()),
3926 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
3927 Proposal::Data { data } => {
3933 if requires_gating(action_kind)
3934 || data.get("revert_of").and_then(Value::as_str).is_some()
3935 {
3936 Ok(())
3937 } else {
3938 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
3939 }
3940 }
3941 }
3942}
3943
3944#[cfg(test)]
3945mod definition_body_tests {
3946 use super::safe_definition_body;
3947
3948 #[test]
3949 fn ordinary_bodies_pass() {
3950 assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
3951 assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
3952 }
3953
3954 #[test]
3955 fn a_body_cannot_close_its_own_block_or_carry_destruction() {
3956 assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
3961 assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
3962 assert!(!safe_definition_body("RECALL facts FORGET abc"));
3964 assert!(!safe_definition_body("recall facts purge older than 1d"));
3965 assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
3966 assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
3968 assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
3970 }
3971}
3972
3973#[cfg(test)]
3974mod plan_edit_tests {
3975 use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
3976 use serde_json::json;
3977
3978 fn plan() -> serde_json::Value {
3979 json!({
3980 "nodes": ["fetch", "review", "post"],
3981 "edges": [
3982 {"src": "fetch", "dst": "review"},
3983 {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
3984 ],
3985 "bindings": {"fetch": "sha256:tool1"},
3986 "retries": {"fetch": 1}
3987 })
3988 }
3989
3990 #[test]
3991 fn the_allowlist_admits_thresholds_and_refuses_topology() {
3992 assert!(plan_edit_allowed("edges.1.cond"));
3993 assert!(plan_edit_allowed("edges.1.max_cycles"));
3994 assert!(plan_edit_allowed("retries.fetch"));
3995 for path in [
3998 "nodes",
3999 "nodes.0",
4000 "edges.0.src",
4001 "edges.0.dst",
4002 "edges",
4003 "bindings.fetch",
4004 "edges.x.cond",
4005 "",
4006 ] {
4007 assert!(!plan_edit_allowed(path), "{path} must not be editable");
4008 }
4009 }
4010
4011 #[test]
4012 fn values_are_type_checked_against_the_field() {
4013 assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
4016 assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
4017 assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
4018 assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
4019 assert!(plan_value_ok("retries.fetch", &json!(3)));
4020 assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
4021 assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
4022 assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
4023 assert!(!plan_value_ok("edges.0.src", &json!("other")));
4024 }
4025
4026 #[test]
4027 fn get_reads_through_arrays_and_objects_and_absence_is_null() {
4028 let p = plan();
4029 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
4030 assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
4031 assert_eq!(plan_get(&p, "retries.review"), json!(null));
4034 assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
4035 assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
4036 }
4037
4038 #[test]
4039 fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
4040 let mut p = plan();
4041 assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
4042 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
4043 assert!(plan_set(&mut p, "retries.review", json!(2)));
4044 assert_eq!(plan_get(&p, "retries.review"), json!(2));
4045 assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
4046 assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
4047 }
4048}
4049
4050#[cfg(test)]
4051mod definition_proposal_tests {
4052 use super::is_definition_statement;
4053
4054 #[test]
4055 fn definition_statements_are_recognized_in_both_spellings() {
4056 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
4057 assert!(is_definition_statement(" define template foo AS { x }"));
4058 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
4059 assert!(!is_definition_statement("ADD fact {}"));
4061 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
4062 assert!(!is_definition_statement("FORGET abc"));
4063 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
4067 }
4068}