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 if let Ok(rows) =
715 sub.grains_of_type(crate::model::grain_type::OBSERVATION, Some(crate::eval::HARNESS_NS), opts)
716 {
717 let mut seeded = 0usize;
718 for g in rows {
719 if seeded >= HARNESS_SEED_CAP || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE {
720 break;
721 }
722 let kind = g.str_field("observation_kind").unwrap_or_default();
723 if !HARNESS_EVIDENCE_KINDS.contains(&kind) {
724 continue;
725 }
726 let before = evidence.len();
727 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
728 if evidence.len() > before {
729 seeded += 1;
730 }
731 }
732 }
733 'notes: for ns in &scan_ns {
743 if let Ok(recent) =
744 sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
745 {
746 for g in recent {
747 if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
748 break 'notes;
749 }
750 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
751 }
752 }
753 }
754 'seed: for gt in [
755 crate::model::grain_type::FACT,
756 crate::model::grain_type::OBSERVATION,
757 ] {
758 for ns in &scan_ns {
759 if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
760 for g in recent {
761 if evidence.len() >= EVIDENCE_CAP {
762 break 'seed;
763 }
764 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
765 }
766 }
767 }
768 }
769 funnel.evidence = evidence.len() as u64;
770 if evidence.is_empty() {
771 return Vec::new(); }
773 let (approved, rejected) = self.llm_history(sub);
778 let base = match self.policy.discover_objective {
779 crate::policy::DiscoverObjective::ReviewQueue => DISCOVER_INSTRUCTIONS,
780 crate::policy::DiscoverObjective::Learner => DISCOVER_LEARNER_INSTRUCTIONS,
781 };
782 let mut instructions = base.to_string();
785 if self.policy.skills.enabled {
786 instructions.push_str(&skill_instructions(self.policy.skills.min_steps));
787 }
788 if self.policy.plans.enabled {
789 instructions.push_str(&plan_instructions(self.policy.plans.min_nodes));
790 }
791 let request = crate::llm::LlmRequest {
792 loop_proto: 1,
793 op: "discover",
794 instructions: &instructions,
795 findings: findings.clone(),
796 evidence: evidence.clone(),
797 rejected,
798 approved,
799 };
800 let Ok(body) = serde_json::to_string(&request) else {
801 return Vec::new();
802 };
803 let raw = match llm.complete(&body) {
804 Ok(r) => r,
805 Err(_) => return Vec::new(), };
807 let caps = sub.capabilities();
811 let mut validated: Vec<ValidatedDraft> = Vec::new();
812 let drafts: Vec<_> = crate::llm::parse_discover(&raw)
813 .recommendations
814 .into_iter()
815 .take(crate::llm::MAX_LLM_DRAFTS)
816 .collect();
817 funnel.proposed = drafts.len() as u64;
818 let id_to_hash: std::collections::BTreeMap<&str, &str> = evidence
819 .iter()
820 .map(|e| (e.id.as_str(), e.hash.as_str()))
821 .collect();
822 for d in drafts {
823 let mut cited: Vec<String> = Vec::new();
824 for c in &d.evidence {
825 if let Some(h) = resolve_citation(c, &bundle, &id_to_hash) {
826 if !cited.contains(&h) {
827 cited.push(h);
828 }
829 }
830 }
831 if cited.is_empty() {
832 funnel.dropped_uncited += 1;
833 continue; }
835 let Ok(target) = TargetRef::parse(&d.target) else {
836 funnel.dropped_target += 1;
837 continue;
838 };
839 let tc = target.target_class();
840 if !matches!(tc, "memory" | "query" | "code") {
846 funnel.dropped_target += 1;
847 continue;
848 }
849 let thin = cited.len() < self.policy.min_evidence as usize;
856 if thin {
857 funnel.advisory_thin_evidence += 1;
858 }
859 let resolved = if thin {
860 None
861 } else {
862 resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps, &self.policy)
863 };
864 if tc == "code" && resolved.is_none() {
869 funnel.dropped_target += 1;
870 continue;
871 }
872 validated.push(ValidatedDraft {
873 draft: d,
874 target_ref: target.as_string(),
875 cited,
876 resolved,
877 });
878 }
879 funnel.cited = validated.len() as u64;
880 if validated.is_empty() {
881 return Vec::new();
882 }
883 let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
890 let outcome_metric = self.outcome_metric_template(sub);
891 self.verify_drafts(&**llm, ground, validated, &evidence, outcome_metric, now_ms, funnel)
892 }
893
894 fn outcome_metric_template<S: OmsSubstrate>(
902 &self,
903 sub: &S,
904 ) -> Option<crate::recommendation::MetricSnapshot> {
905 let e = self.policy.outcome_evalset.as_ref()?;
906 let run = crate::eval::newest_eval_run(sub, &e.hash, None).ok().flatten()?;
907 let baseline = crate::eval::run_value(&run, &e.field)?;
908 let schedule = e.schedule();
912 let ms_only: Vec<i64> = schedule.iter().filter_map(Checkpoint::as_ms).collect();
913 let all_ms = ms_only.len() == schedule.len();
914 let horizons = if ms_only.is_empty() { vec![86_400_000] } else { ms_only };
915 Some(crate::recommendation::MetricSnapshot {
916 metric: format!("evalset:{}:{}", e.hash, e.field),
917 baseline,
918 unit: e.field.clone(),
919 n: run.total(),
920 window: "per-run".into(),
921 subject: None,
922 namespace: None,
923 relation: None,
924 query: format!(
925 "RECALL facts WHERE subject = \"evalset:{}\" AND relation = \"mg:eval_run\"",
926 e.hash
927 ),
928 review_after_ms: horizons[0],
929 horizons_ms: if all_ms { horizons } else { Vec::new() },
930 checkpoints: if all_ms { Vec::new() } else { schedule },
931 higher_is_better: e.higher_is_better,
932 })
933 }
934
935 #[allow(clippy::too_many_arguments)]
942 fn verify_drafts(
943 &self,
944 llm: &dyn crate::llm::LlmBackend,
945 ground: &dyn crate::llm::LlmBackend,
946 validated: Vec<ValidatedDraft>,
947 evidence: &[crate::llm::EvidenceItem],
948 outcome_metric: Option<crate::recommendation::MetricSnapshot>,
949 now_ms: i64,
950 funnel: &mut LlmFunnel,
951 ) -> Vec<Recommendation> {
952 use crate::llm::*;
953 let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
954 evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
955 let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
956 cited
957 .iter()
958 .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
959 .collect()
960 };
961
962 let claims: Vec<GroundItem> = validated
966 .iter()
967 .enumerate()
968 .map(|(i, v)| GroundItem {
969 id: i,
970 claim: claim_text(&v.draft, v.resolved.as_ref()),
971 evidence: ev_for(&v.cited),
972 })
973 .collect();
974 let ground_req = GroundRequest {
975 loop_proto: 1,
976 op: "ground",
977 instructions: GROUND_INSTRUCTIONS,
978 claims,
979 };
980 let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
986 .ok()
987 .and_then(|b| ground.complete(&b).ok())
988 {
989 Some(raw) => {
990 let parsed = parse_ground(&raw);
991 funnel.ground_verdicts = parsed.results.len() as u64;
992 parsed
993 .results
994 .into_iter()
995 .filter(|r| r.supported)
996 .map(|r| r.id)
997 .collect()
998 }
999 None => {
1000 funnel.ground_call_failed = true;
1001 return Vec::new();
1002 }
1003 };
1004 funnel.grounded = grounded.len() as u64;
1005 if grounded.is_empty() {
1006 return Vec::new();
1007 }
1008
1009 let items: Vec<VerifyItem> = validated
1015 .iter()
1016 .enumerate()
1017 .filter(|(i, _)| grounded.contains(i))
1018 .map(|(i, v)| VerifyItem {
1019 id: i,
1020 summary: claim_text(&v.draft, v.resolved.as_ref()),
1022 target: v.target_ref.clone(),
1023 evidence: ev_for(&v.cited),
1024 })
1025 .collect();
1026 let verify_req = VerifyRequest {
1027 loop_proto: 1,
1028 op: "verify",
1029 instructions: VERIFY_INSTRUCTIONS,
1030 findings: items,
1031 };
1032 let verdicts: std::collections::BTreeMap<usize, f64> =
1033 match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
1034 Some(raw) => parse_verify(&raw)
1035 .results
1036 .into_iter()
1037 .filter(|r| r.keep)
1038 .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
1039 .collect(),
1040 None => return Vec::new(),
1041 };
1042
1043 funnel.kept = verdicts.len() as u64;
1044 let mut out = Vec::new();
1048 for (i, v) in validated.into_iter().enumerate() {
1049 if let Some(&conf) = verdicts.get(&i) {
1050 if conf >= MIN_LLM_CONFIDENCE {
1051 let mut rec = stamp_llm(
1052 llm.model(),
1053 &v.draft,
1054 v.target_ref,
1055 v.cited,
1056 v.resolved,
1057 conf,
1058 now_ms,
1059 );
1060 if rec.rollbackable {
1064 rec.metric = outcome_metric.clone();
1065 }
1066 out.push(rec);
1067 }
1068 }
1069 }
1070 funnel.stored = out.len() as u64;
1071 out
1072 }
1073
1074 fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
1079 const MAX: usize = 20;
1080 let Ok(mut recs) = self.recommendations(sub, None) else {
1081 return (Vec::new(), Vec::new());
1082 };
1083 recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
1084 recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
1085 let mut approved = Vec::new();
1086 let mut rejected = Vec::new();
1087 for r in &recs {
1088 match r.status {
1089 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
1090 if approved.len() < MAX =>
1091 {
1092 approved.push(r.summary.render());
1093 }
1094 RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
1095 _ => {}
1096 }
1097 }
1098 (approved, rejected)
1099 }
1100
1101 fn enrich(&self, survivors: &mut [Recommendation]) {
1106 let Some(llm) = &self.llm else {
1107 return;
1108 };
1109 if survivors.is_empty() {
1110 return;
1111 }
1112 let findings: Vec<crate::llm::FindingBrief> = survivors
1113 .iter()
1114 .map(|r| crate::llm::FindingBrief {
1115 analyzer: r.analyzer.clone(),
1116 summary: r.summary.render(),
1117 target: r.target_ref.clone(),
1118 severity: r.severity.as_str().to_string(),
1119 })
1120 .collect();
1121 let request = crate::llm::LlmRequest {
1122 loop_proto: 1,
1123 op: "enrich",
1124 instructions: ENRICH_INSTRUCTIONS,
1125 findings,
1126 evidence: Vec::new(),
1127 rejected: Vec::new(),
1128 approved: Vec::new(),
1129 };
1130 let Ok(body) = serde_json::to_string(&request) else {
1131 return;
1132 };
1133 let raw = match llm.complete(&body) {
1134 Ok(r) => r,
1135 Err(_) => return,
1136 };
1137 for note in crate::llm::parse_enrich(&raw).notes {
1138 if note.guidance.trim().is_empty() {
1139 continue;
1140 }
1141 if let Some(r) = survivors
1142 .iter_mut()
1143 .find(|r| r.target_ref == note.target && r.guidance.is_none())
1144 {
1145 r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
1146 }
1147 }
1148 }
1149
1150 fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
1159 if !rec.origin.auto_apply_eligible() || rec.destructive {
1160 return false;
1161 }
1162 let manifest_ok = self
1166 .analyzers
1167 .iter()
1168 .map(|a| a.manifest())
1169 .find(|m| m.id == rec.analyzer)
1170 .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
1171 if !manifest_ok {
1172 return false;
1173 }
1174 let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1175 return false;
1176 };
1177 let family = crate::manifest::analyzer_family(&rec.analyzer);
1178 if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1179 return false;
1180 }
1181 match &rec.proposal {
1186 Proposal::Cal { cal } => cal
1187 .lines()
1188 .map(str::trim)
1189 .filter(|l| !l.is_empty())
1190 .all(|l| supersede_is_value_identical(sub, l)),
1191 _ => false,
1192 }
1193 }
1194
1195 fn auto_apply<S: OmsSubstrate>(
1198 &self,
1199 sub: &mut S,
1200 p: &mut LoopPersisted,
1201 rec: &Recommendation,
1202 now_ms: i64,
1203 ) -> Result<()> {
1204 let mut created = Vec::new();
1205 if let Proposal::Cal { cal } = &rec.proposal {
1206 if cal.lines().map(str::trim).any(is_definition_statement) {
1211 return Err(Error::InvalidProposal(
1212 "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1213 auto-applied: it changes what every future context contains, so it \
1214 requires a human APPROVE + APPLY with BECAUSE"
1215 .into(),
1216 ));
1217 }
1218 for r in sub.execute_cal(cal)? {
1219 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1220 created.push(h.to_string());
1221 }
1222 }
1223 }
1224 let applied = AppliedRecord {
1225 applied_at_ms: now_ms,
1226 target_ref: rec.target_ref.clone(),
1227 rollbackable: rec.rollbackable,
1228 created_hashes: created,
1229 inverse_cal: None,
1230 metric: rec.metric.clone(),
1231 };
1232 let prev = p.audit_heads.get(&rec.hash).cloned();
1233 let audit = AuditRecord {
1234 rec_hash: rec.hash.clone(),
1235 from: Some(RecStatus::Pending),
1236 to: RecStatus::Applied,
1237 actor: "policy:auto".into(),
1238 observer_type: ObserverType::Policy,
1239 because: "auto-applied per host policy".into(),
1240 previous_audit_hash: prev,
1241 gating: None,
1242 at_ms: now_ms,
1243 };
1244 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1245 p.audit_heads.insert(rec.hash.clone(), audit_hash);
1246 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1247 p.applied.insert(rec.hash.clone(), applied);
1248 Ok(())
1249 }
1250
1251 #[allow(clippy::too_many_arguments)]
1254 pub fn review<S: OmsSubstrate>(
1255 &self,
1256 sub: &mut S,
1257 rec_hash: &str,
1258 decision: Decision,
1259 actor: &str,
1260 observer: ObserverType,
1261 scopes: &ScopeSet,
1262 because: &str,
1263 now_ms: i64,
1264 ) -> Result<()> {
1265 if !scopes.has(Scope::Review) {
1266 return Err(Error::ScopeDenied("review".into()));
1267 }
1268 let because = validate_because(because)?;
1269 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1270 let status = *p
1271 .status_index
1272 .get(rec_hash)
1273 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1274 let to = match decision {
1275 Decision::Approve => RecStatus::Approved,
1276 Decision::Reject => RecStatus::Rejected,
1277 };
1278 if !status.can_transition_to(to, false) {
1279 return Err(Error::LifecycleViolation(format!(
1280 "{} -> {}",
1281 status.as_str(),
1282 to.as_str()
1283 )));
1284 }
1285 if to == RecStatus::Approved {
1286 if let Some(creator) = p.creators.get(rec_hash) {
1287 if creator == actor {
1288 return Err(Error::SelfApproval(format!(
1289 "{actor} created this recommendation"
1290 )));
1291 }
1292 }
1293 if let Some(trigger) = p.co_creators.get(rec_hash) {
1294 if trigger == actor {
1295 return Err(Error::SelfApproval(format!(
1296 "{actor} triggered the run that authored this recommendation"
1297 )));
1298 }
1299 }
1300 }
1301 let prev = p.audit_heads.get(rec_hash).cloned();
1302 let audit = AuditRecord {
1303 rec_hash: rec_hash.into(),
1304 from: Some(status),
1305 to,
1306 actor: actor.into(),
1307 observer_type: observer,
1308 because,
1309 previous_audit_hash: prev,
1310 gating: None,
1311 at_ms: now_ms,
1312 };
1313 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1314 p.audit_heads.insert(rec_hash.into(), audit_hash);
1315 p.status_index.insert(rec_hash.into(), to);
1316 if to == RecStatus::Rejected {
1317 if let Ok(rec) = load_rec(sub, rec_hash) {
1318 strike_cooldown(&mut p, rec.dedup_key, now_ms);
1319 }
1320 }
1321 sub.store_state(&p.to_value()?)?;
1322 Ok(())
1323 }
1324
1325 pub fn preflight_apply<S: OmsSubstrate>(
1342 &self,
1343 sub: &S,
1344 rec_hash: &str,
1345 scopes: &ScopeSet,
1346 allow_destructive: bool,
1347 has_gating: bool,
1348 ) -> Result<()> {
1349 if !scopes.has(Scope::Apply) {
1350 return Err(Error::ScopeDenied("apply".into()));
1351 }
1352 let rec = load_rec(sub, rec_hash)?;
1353 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1354 return Err(Error::DestructiveGated(
1355 "destructive apply requires admin scope + allow_destructive".into(),
1356 ));
1357 }
1358 ensure_executable(rec.action_kind, &rec.proposal)?;
1359 if requires_gating(rec.action_kind) && !has_gating {
1360 return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1361 }
1362 Ok(())
1363 }
1364
1365 #[allow(clippy::too_many_arguments)]
1369 pub fn apply<S: OmsSubstrate>(
1370 &self,
1371 sub: &mut S,
1372 rec_hash: &str,
1373 actor: &str,
1374 observer: ObserverType,
1375 scopes: &ScopeSet,
1376 because: &str,
1377 allow_destructive: bool,
1378 now_ms: i64,
1379 ) -> Result<AppliedRecord> {
1380 self.apply_inner(
1381 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1382 )
1383 }
1384
1385 pub fn gating_evidence<S: OmsSubstrate>(
1392 &self,
1393 sub: &S,
1394 rec_hash: &str,
1395 run_id: &str,
1396 ) -> Result<crate::recommendation::GatingEvidence> {
1397 let rec = self
1398 .recommendations(sub, None)?
1399 .into_iter()
1400 .find(|r| r.hash == rec_hash)
1401 .ok_or_else(|| {
1402 Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1403 })?;
1404 let pin = rec.evalset_hash.ok_or_else(|| {
1405 Error::InvalidProposal(
1406 "this recommendation pins no evalset — a gating run applies only \
1407 to code and adapter revisions"
1408 .into(),
1409 )
1410 })?;
1411 match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1416 Some(run) => Ok(crate::recommendation::GatingEvidence {
1417 evalset_hash: pin,
1418 run_id: run.run_id,
1419 passed: run.passed,
1420 failed: run.failed,
1421 }),
1422 None => Err(Error::InvalidProposal(format!(
1423 "no recorded gate run '{run_id}' for evalset {pin} — run \
1424 `areev eval run --evalset {pin} ...` first"
1425 ))),
1426 }
1427 }
1428
1429 #[allow(clippy::too_many_arguments)]
1433 pub fn apply_gated<S: OmsSubstrate>(
1434 &self,
1435 sub: &mut S,
1436 rec_hash: &str,
1437 actor: &str,
1438 observer: ObserverType,
1439 scopes: &ScopeSet,
1440 because: &str,
1441 allow_destructive: bool,
1442 gating: &crate::recommendation::GatingEvidence,
1443 now_ms: i64,
1444 ) -> Result<AppliedRecord> {
1445 self.apply_inner(
1446 sub,
1447 rec_hash,
1448 actor,
1449 observer,
1450 scopes,
1451 because,
1452 allow_destructive,
1453 Some(gating),
1454 now_ms,
1455 )
1456 }
1457
1458 #[allow(clippy::too_many_arguments)]
1459 fn apply_inner<S: OmsSubstrate>(
1460 &self,
1461 sub: &mut S,
1462 rec_hash: &str,
1463 actor: &str,
1464 observer: ObserverType,
1465 scopes: &ScopeSet,
1466 because: &str,
1467 allow_destructive: bool,
1468 gating: Option<&crate::recommendation::GatingEvidence>,
1469 now_ms: i64,
1470 ) -> Result<AppliedRecord> {
1471 if !scopes.has(Scope::Apply) {
1472 return Err(Error::ScopeDenied("apply".into()));
1473 }
1474 let because = validate_because(because)?;
1475 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1476 let status = *p
1477 .status_index
1478 .get(rec_hash)
1479 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1480 if !status.can_transition_to(RecStatus::Applied, false) {
1481 return Err(Error::LifecycleViolation(format!(
1482 "{} -> applied (approve first)",
1483 status.as_str()
1484 )));
1485 }
1486 let rec = load_rec(sub, rec_hash)?;
1487 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1488 return Err(Error::DestructiveGated(
1489 "destructive apply requires admin scope + allow_destructive".into(),
1490 ));
1491 }
1492 if requires_gating(rec.action_kind) {
1497 let g = gating
1498 .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1499 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1500 if g.evalset_hash != pin {
1501 return Err(Error::InvalidProposal(format!(
1502 "gating ran evalset {} but the recommendation is pinned \
1503 to {pin} (Rule E1)",
1504 g.evalset_hash
1505 )));
1506 }
1507 match sub.grain(pin)? {
1508 Some(evalset) if evalset.is_live() => {}
1509 Some(_) => {
1510 return Err(Error::InvalidProposal(
1511 "the pinned evalset was superseded after gating — \
1512 the recommendation must re-gate (Rule E1)"
1513 .into(),
1514 ))
1515 }
1516 None => {
1517 return Err(Error::InvalidProposal(format!(
1518 "pinned evalset {pin} not found in the substrate"
1519 )))
1520 }
1521 }
1522 if g.failed > 0 {
1523 return Err(Error::InvalidProposal(format!(
1524 "the gating run failed {}/{} cases — a failing gate \
1525 admits nothing",
1526 g.failed,
1527 g.passed + g.failed
1528 )));
1529 }
1530 }
1531
1532 let mut created = Vec::new();
1534 let mut inverse_cal: Option<String> = None;
1538 match &rec.proposal {
1539 Proposal::Cal { cal } => {
1540 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1541 if !is_definition_statement(line) {
1542 continue;
1543 }
1544 match sub.definition_inverse(line)? {
1545 Some(inv) => inverse_cal = Some(inv),
1546 None => {
1547 return Err(Error::InvalidProposal(format!(
1548 "this substrate cannot record a rollback inverse for {line:?}; \
1549 a definition rewrite that ROLLBACK could not undo is refused \
1550 rather than applied"
1551 )))
1552 }
1553 }
1554 }
1555 let rows = sub.execute_cal(cal)?;
1556 for r in rows {
1557 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1558 created.push(h.to_string());
1559 }
1560 }
1561 }
1562 Proposal::Edit { .. } => {
1565 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1566 }
1567 Proposal::Data { data } if requires_gating(rec.action_kind) => {
1575 let relation = if rec.action_kind == ActionKind::AdapterRevision {
1576 "mg:adapter_promotion"
1577 } else {
1578 "mg:code_promotion"
1579 };
1580 let mut promoted = data.clone();
1588 if let Some(Value::String(src)) = promoted.remove("source") {
1589 let address = sub.put_blob(src.as_bytes())?;
1590 promoted.insert("code_address".into(), Value::from(address));
1591 }
1592 let mut spec = crate::substrate::GrainSpec::new(
1593 crate::model::grain_type::FACT,
1594 LOOP_NS,
1595 )
1596 .with_field("subject", rec.target_ref.clone())
1597 .with_field("relation", relation)
1598 .with_field(
1599 "object",
1600 serde_json::to_string(&promoted)
1601 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1602 )
1603 .with_field("rec_hash", rec_hash.to_string());
1604 if let Some(g) = gating {
1605 spec = spec
1606 .with_field("gating_evalset", g.evalset_hash.clone())
1607 .with_field("gating_run_id", g.run_id.clone());
1608 }
1609 created.push(sub.put_grain(&spec)?);
1610 }
1611 Proposal::Data { data } => {
1612 let revert_of = data
1617 .get("revert_of")
1618 .and_then(Value::as_str)
1619 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1620 self.rollback(
1621 sub,
1622 revert_of,
1623 actor,
1624 observer,
1625 scopes,
1626 &because,
1627 now_ms,
1628 )?;
1629 p = LoopPersisted::from_value(sub.load_state()?)?;
1632 if let Ok(reverted) = load_rec(sub, revert_of) {
1642 strike_cooldown(&mut p, reverted.dedup_key, now_ms);
1643 }
1644 }
1645 }
1646
1647 let applied = AppliedRecord {
1648 applied_at_ms: now_ms,
1649 target_ref: rec.target_ref.clone(),
1650 rollbackable: rec.rollbackable,
1651 created_hashes: created,
1652 inverse_cal,
1653 metric: rec.metric.clone(),
1654 };
1655 let prev = p.audit_heads.get(rec_hash).cloned();
1656 let audit = AuditRecord {
1657 rec_hash: rec_hash.into(),
1658 from: Some(status),
1659 to: RecStatus::Applied,
1660 actor: actor.into(),
1661 observer_type: observer,
1662 because,
1663 previous_audit_hash: prev,
1664 gating: gating.cloned(),
1665 at_ms: now_ms,
1666 };
1667 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1668 p.audit_heads.insert(rec_hash.into(), audit_hash);
1669 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1670 p.applied.insert(rec_hash.into(), applied.clone());
1671 sub.store_state(&p.to_value()?)?;
1672 Ok(applied)
1673 }
1674
1675 #[allow(clippy::too_many_arguments)]
1678 pub fn rollback<S: OmsSubstrate>(
1679 &self,
1680 sub: &mut S,
1681 rec_hash: &str,
1682 actor: &str,
1683 observer: ObserverType,
1684 scopes: &ScopeSet,
1685 because: &str,
1686 now_ms: i64,
1687 ) -> Result<()> {
1688 if !scopes.has(Scope::Apply) {
1689 return Err(Error::ScopeDenied("apply".into()));
1690 }
1691 let because = validate_because(because)?;
1692 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1693 let status = *p
1694 .status_index
1695 .get(rec_hash)
1696 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1697 if !status.can_transition_to(RecStatus::RolledBack, false) {
1698 return Err(Error::LifecycleViolation(format!(
1699 "{} -> rolled_back",
1700 status.as_str()
1701 )));
1702 }
1703 let applied = p
1704 .applied
1705 .get(rec_hash)
1706 .cloned()
1707 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1708 if !applied.rollbackable {
1709 return Err(Error::LifecycleViolation(
1710 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1711 ));
1712 }
1713 for h in &applied.created_hashes {
1714 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1715 }
1716 if let Some(inverse) = &applied.inverse_cal {
1723 sub.execute_cal(inverse)?;
1724 }
1725 let prev = p.audit_heads.get(rec_hash).cloned();
1726 let audit = AuditRecord {
1727 rec_hash: rec_hash.into(),
1728 from: Some(status),
1729 to: RecStatus::RolledBack,
1730 actor: actor.into(),
1731 observer_type: observer,
1732 because,
1733 previous_audit_hash: prev,
1734 gating: None,
1735 at_ms: now_ms,
1736 };
1737 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1738 p.audit_heads.insert(rec_hash.into(), audit_hash);
1739 p.status_index
1740 .insert(rec_hash.into(), RecStatus::RolledBack);
1741 sub.store_state(&p.to_value()?)?;
1742 Ok(())
1743 }
1744
1745 pub fn recommendations<S: OmsSubstrate>(
1750 &self,
1751 sub: &S,
1752 status_filter: Option<RecStatus>,
1753 ) -> Result<Vec<Recommendation>> {
1754 let p = LoopPersisted::from_value(sub.load_state()?)?;
1755 let grains = sub.grains_of_type(
1756 crate::model::grain_type::RECOMMENDATION,
1757 Some(LOOP_NS),
1758 ReadOpts {
1759 live_only: false,
1760 since_ms: None,
1761 },
1762 )?;
1763 let mut out = Vec::new();
1764 for g in grains {
1765 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1766 rec.status = p
1767 .status_index
1768 .get(&g.hash)
1769 .copied()
1770 .unwrap_or(RecStatus::Pending);
1771 if let Some(f) = status_filter {
1772 if rec.status != f {
1773 continue;
1774 }
1775 }
1776 out.push(rec);
1777 }
1778 out.sort_by(|a, b| {
1787 b.severity
1788 .cmp(&a.severity)
1789 .then(a.created_at_ms.cmp(&b.created_at_ms))
1790 .then(a.dedup_key.cmp(&b.dedup_key))
1791 .then(a.hash.cmp(&b.hash))
1792 });
1793 Ok(out)
1794 }
1795
1796 pub fn analyzer_settings<S: OmsSubstrate>(
1799 &self,
1800 sub: &S,
1801 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1802 let p = LoopPersisted::from_value(sub.load_state()?)?;
1803 Ok(self
1804 .analyzers
1805 .iter()
1806 .map(|a| {
1807 let m = a.manifest();
1808 let cfg = p.config.get(&m.id);
1809 crate::config::AnalyzerSetting {
1810 id: m.id.clone(),
1811 title: m.title.clone(),
1812 description: m.description.clone(),
1813 tier: format!("{:?}", m.tier),
1814 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1815 default_on: m.default_on,
1816 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1817 severity_floor: cfg
1818 .and_then(|c| c.severity_floor)
1819 .map(|s| s.as_str().to_string()),
1820 }
1821 })
1822 .collect())
1823 }
1824
1825 pub fn set_analyzer_config<S: OmsSubstrate>(
1831 &self,
1832 sub: &mut S,
1833 analyzer_id: &str,
1834 update: crate::config::AnalyzerConfigUpdate,
1835 scopes: &ScopeSet,
1836 ) -> Result<crate::config::AnalyzerConfig> {
1837 if !scopes.has(Scope::Admin) {
1838 return Err(Error::ScopeDenied("admin".into()));
1839 }
1840 let manifest = self
1841 .analyzers
1842 .iter()
1843 .map(|a| a.manifest())
1844 .find(|m| m.id == analyzer_id)
1845 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1846 if let Some(params) = &update.params {
1848 manifest.resolve_params(params)?;
1849 }
1850 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1851 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1852 if let Some(enabled) = update.enabled {
1853 cfg.enabled = Some(enabled);
1854 }
1855 if update.clear_floor {
1856 cfg.severity_floor = None;
1857 } else if let Some(floor) = update.severity_floor {
1858 cfg.severity_floor = Some(floor);
1859 }
1860 if let Some(params) = update.params {
1861 cfg.params = params;
1862 }
1863 if let Some(ns) = update.namespaces {
1864 cfg.namespaces = ns;
1865 }
1866 let stored = cfg.clone();
1867 sub.store_state(&p.to_value()?)?;
1868 Ok(stored)
1869 }
1870
1871 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1874 let p = LoopPersisted::from_value(sub.load_state()?)?;
1875 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1876 out.sort_by(|a, b| {
1880 a.measured_at_ms
1881 .cmp(&b.measured_at_ms)
1882 .then(a.horizon_ms.cmp(&b.horizon_ms))
1883 .then(a.metric.cmp(&b.metric))
1884 .then(a.rec_hash.cmp(&b.rec_hash))
1885 });
1886 Ok(out)
1887 }
1888
1889 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1893 let p = LoopPersisted::from_value(sub.load_state()?)?;
1894 let new = count_new(sub, p.state.watermark_ms)?;
1895 let (grains_since_run, error_events_since_run) = (new.grains, new.error_events);
1896 let recs = self.recommendations(sub, None)?;
1897 let mut pending = 0;
1898 let mut applied = 0;
1899 for r in &recs {
1900 match r.status {
1901 RecStatus::Pending => pending += 1,
1902 RecStatus::Applied => applied += 1,
1903 _ => {}
1904 }
1905 }
1906 let stale = match p.state.last_run_ms {
1908 None => true,
1909 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1910 };
1911 Ok(Health {
1912 last_run_ms: p.state.last_run_ms,
1913 grains_since_run,
1914 error_events_since_run,
1915 pending,
1916 applied,
1917 total: recs.len() as u64,
1918 stale,
1919 })
1920 }
1921
1922 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1927 let recs = self.recommendations(sub, None)?;
1928 let mut m = LlmMetrics::default();
1929 for r in &recs {
1930 if !matches!(r.origin, Origin::Llm { .. }) {
1931 continue;
1932 }
1933 m.proposed += 1;
1934 match r.status {
1935 RecStatus::Pending => m.pending += 1,
1936 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1937 RecStatus::Rejected => m.rejected += 1,
1938 RecStatus::Expired => {}
1939 }
1940 }
1941 let decided = m.approved + m.rejected;
1942 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1943 Ok(m)
1944 }
1945}
1946
1947#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1949pub struct Health {
1950 #[serde(skip_serializing_if = "Option::is_none")]
1951 pub last_run_ms: Option<i64>,
1952 pub grains_since_run: u64,
1953 pub error_events_since_run: u64,
1954 pub pending: u64,
1955 pub applied: u64,
1956 pub total: u64,
1957 pub stale: bool,
1960}
1961
1962#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1964pub struct LlmMetrics {
1965 pub proposed: u64,
1968 pub pending: u64,
1969 pub approved: u64,
1971 pub rejected: u64,
1972 #[serde(skip_serializing_if = "Option::is_none")]
1974 pub approval_rate: Option<f64>,
1975}
1976
1977fn measure_outcomes<S: OmsSubstrate>(
1984 sub: &S,
1985 p: &mut LoopPersisted,
1986 now_ms: i64,
1987) -> Result<Vec<OutcomeInput>> {
1988 let mut due: Vec<(String, crate::config::AppliedRecord, Checkpoint)> = Vec::new();
1992 for (h, a) in &p.applied {
1993 if p.status_index.get(h) != Some(&RecStatus::Applied) {
1994 continue;
1995 }
1996 let Some(metric) = &a.metric else { continue };
1997 let done = p.measured.get(h).cloned().unwrap_or_default();
1998 for cp in metric.schedule() {
1999 if !done.contains(&cp) && checkpoint_due(sub, metric, a.applied_at_ms, cp, now_ms)? {
2000 due.push((h.clone(), a.clone(), cp));
2001 }
2002 }
2003 }
2004
2005 let mut out = Vec::new();
2006 for (rec_hash, applied, checkpoint) in due {
2007 let metric = applied.metric.as_ref().unwrap();
2008 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
2009 continue; };
2011 let baseline = baseline_at_apply(sub, metric, applied.applied_at_ms)?;
2012 let regressed = crate::recommendation::is_regression(
2013 baseline,
2014 current,
2015 metric.higher_is_better,
2016 );
2017 p.outcomes.entry(rec_hash.clone()).or_default().push(
2018 crate::recommendation::OutcomeResult {
2019 rec_hash: rec_hash.clone(),
2020 metric: metric.metric.clone(),
2021 baseline,
2022 current,
2023 verdict: if regressed { "regressed" } else { "held" }.into(),
2024 horizon_ms: checkpoint.as_ms().unwrap_or(0),
2027 checkpoint: checkpoint.as_ms().is_none().then_some(checkpoint),
2028 measured_at_ms: now_ms,
2029 },
2030 );
2031 p.measured.entry(rec_hash.clone()).or_default().push(checkpoint);
2032 if regressed {
2033 out.push(OutcomeInput {
2034 rec_hash,
2035 target_ref: applied.target_ref.clone(),
2036 metric: metric.metric.clone(),
2037 baseline,
2038 current,
2039 unit: metric.unit.clone(),
2040 higher_is_better: metric.higher_is_better,
2041 });
2042 }
2043 }
2044 Ok(out)
2045}
2046
2047fn checkpoint_due<S: SubstrateRead>(
2056 sub: &S,
2057 metric: &crate::recommendation::MetricSnapshot,
2058 applied_at_ms: i64,
2059 cp: Checkpoint,
2060 now_ms: i64,
2061) -> Result<bool> {
2062 Ok(match cp {
2063 Checkpoint::AfterMs(ms) => now_ms - applied_at_ms >= ms,
2064 Checkpoint::AfterRuns(n) => match crate::eval::parse_evalset_metric(&metric.metric) {
2065 Some((evalset, _)) => {
2066 crate::eval::eval_runs(sub, evalset, Some(applied_at_ms + 1))?.len() >= n as usize
2067 }
2068 None => false,
2069 },
2070 Checkpoint::AfterGrains(n) => count_new(sub, Some(applied_at_ms))?.grains >= n as u64,
2071 })
2072}
2073
2074fn baseline_at_apply<S: SubstrateRead>(
2086 sub: &S,
2087 metric: &crate::recommendation::MetricSnapshot,
2088 applied_at_ms: i64,
2089) -> Result<f64> {
2090 if let Some((evalset, field)) = crate::eval::parse_evalset_metric(&metric.metric) {
2091 if let Some(run) = crate::eval::newest_eval_run_before(sub, evalset, applied_at_ms)? {
2092 if let Some(v) = crate::eval::run_value(&run, field) {
2093 return Ok(v);
2094 }
2095 }
2096 }
2097 Ok(metric.baseline)
2098}
2099
2100pub(crate) fn measure_metric<S: SubstrateRead>(
2102 sub: &S,
2103 metric: &crate::recommendation::MetricSnapshot,
2104 since_ms: i64,
2105) -> Result<Option<f64>> {
2106 match metric.metric.as_str() {
2107 "tool_error_recurrence" => {
2112 let Some(tool) = &metric.subject else { return Ok(None) };
2113 let tools = sub.grains_of_type(
2114 crate::model::grain_type::TOOL,
2115 None,
2116 ReadOpts { live_only: true, since_ms: Some(since_ms) },
2117 )?;
2118 let n = tools
2119 .iter()
2120 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
2121 .filter(|t| {
2122 metric.relation.as_deref().is_none_or(|sig| {
2125 crate::analyzers::tool_failure::normalize_signature(
2126 t.tool_content().unwrap_or(""),
2127 ) == sig
2128 })
2129 })
2130 .count();
2131 Ok(Some(n as f64))
2132 }
2133 "contradiction_recurrence" => {
2137 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
2138 return Ok(None);
2139 };
2140 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
2141 let distinct: BTreeSet<String> = facts
2142 .iter()
2143 .filter(|f| {
2144 f.fact_relation()
2145 .is_some_and(|r| normalize_ident(r) == *relation)
2146 })
2147 .filter_map(|f| f.fact_object().map(normalize_ident))
2148 .collect();
2149 Ok(Some(distinct.len().saturating_sub(1) as f64))
2150 }
2151 m if m.starts_with("evalset:") => {
2164 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
2165 return Ok(None);
2166 };
2167 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
2168 return Ok(None);
2169 };
2170 Ok(crate::eval::run_value(&run, field))
2171 }
2172 _ => Ok(None),
2173 }
2174}
2175
2176fn scoped_live_facts<S: SubstrateRead>(
2179 sub: &S,
2180 namespace: Option<&str>,
2181 subject: &str,
2182) -> Result<Vec<GrainRecord>> {
2183 let facts = sub.grains_of_type(
2184 crate::model::grain_type::FACT,
2185 None,
2186 ReadOpts { live_only: true, since_ms: None },
2187 )?;
2188 Ok(facts
2189 .into_iter()
2190 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
2191 .filter(|f| {
2192 f.fact_subject()
2193 .is_some_and(|s| normalize_ident(s) == subject)
2194 })
2195 .collect())
2196}
2197
2198fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
2210 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
2211 return false;
2212 };
2213 if fields.is_empty() {
2214 return false;
2215 }
2216 let Ok(Some(grain)) = sub.grain(&target) else {
2217 return false;
2218 };
2219 if grain.valid_to_ms.is_some() {
2229 return false;
2230 }
2231 fields.iter().all(|(k, v)| {
2233 if k == "namespace" {
2234 return v
2235 .as_str()
2236 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2237 }
2238 match (v, grain.fields.get(k)) {
2239 (Value::String(a), Some(Value::String(b))) => {
2240 normalize_ident(a) == normalize_ident(b)
2241 }
2242 (a, Some(b)) => a == b,
2243 (_, None) => false,
2244 }
2245 })
2246}
2247
2248const EVIDENCE_CAP: usize = 64;
2258const CITED_SEED_CAP: usize = 24;
2259const TOOL_SEED_CAP: usize = 16;
2260const NOTE_SEED_CAP: usize = 8;
2265const HARNESS_EVIDENCE_KINDS: &[&str] = &["fold_summary"];
2284const HARNESS_SEED_CAP: usize = 6;
2285const LENS_RESERVE: usize = 24;
2286
2287const MIN_LLM_CONFIDENCE: f64 = 0.75;
2290
2291macro_rules! discover_instructions {
2301 ($scoring:literal) => {
2302 concat!(
2303 "You review an agent's memory for quality. \
2304Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2305findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2306stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2307recurring cost or hand-off the agent's own setup could remove). \
2308The deterministic findings already cover what the ERROR TEXT says; restating \
2309one of them earns nothing. The evidence may also contain OUTCOME records — a \
2310run's observable shape together with whether it was accepted or rejected. A \
2311problem that raised no error at all is exactly the kind the deterministic \
2312checks cannot see, so compare the rejected outcomes against the accepted \
2313ones: a feature they share and the accepted ones lack is a candidate rule. \
2314Require at least two rejected outcomes before proposing one — a single \
2315rejection is an anecdote, not a pattern. ",
2316 $scoring,
2317 " The 'approved' and 'rejected' lists, when \
2318present, show findings this reviewer recently accepted or rejected — prefer the \
2319kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2320or more evidence items from the bundle by their 'id' (or 'hash'), name a \
2321'target', and include your confidence 0.0-1.0. Return JSON: \
2322{\"recommendations\":[{\"summary\":\"...\",\
2323\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<id>\"],\
2324\"confidence\":0.0,\"proposal\":{...}}]}. \
2325OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2326you are not asking to change anything. Include it ONLY when the evidence \
2327supports a specific change, choosing exactly one kind: \
2328(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2329\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
2330an action the agent itself takes on the next occasion. Either ADD an action \
2331it is failing to take ('Record the vendor name and the amount on every \
2332invoice, not just the date') or ORDER one that goes wrong ('Refund a \
2333subscription before cancelling it; refunds on cancelled subscriptions are \
2334refused'). Name the action, not a check on it: 'validate', 'verify' and \
2335'ensure ... is correct' describe a review step the agent has no way to \
2336perform, and such a rule changes nothing even once applied. \
2337(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2338entity target — a durable fact the agent keeps having to be told (an alias, a \
2339settled default, a preference). 'relation' is a short identifier (letters, \
2340digits, _ - . :), not a sentence. \
2341(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2342\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2343assembles the agent's context, when the evidence shows it retrieves the wrong \
2344things. Give the FULL new body; it replaces the old one. \
2345(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2346with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2347workflow. Only these paths are editable: 'edges.<i>.cond', \
2348'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2349holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2350(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2351— full replacement source for that tool. It is applied only after a recorded \
2352evaluation run passes, so propose one only when the evidence shows the current \
2353code is the defect. \
2354The subject of a fact, the name of a query, the plan hash and the tool name \
2355all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2356becomes a change a human reviewer may apply, so it must be fully supported by \
2357the cited evidence. Propose nothing you cannot ground in the evidence."
2358 )
2359 };
2360}
2361
2362const DISCOVER_INSTRUCTIONS: &str = discover_instructions!(
2364 "SCORING: propose a finding ONLY if you \
2365are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2366useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2367earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2368when there is nothing worth flagging."
2369);
2370
2371const DISCOVER_LEARNER_INSTRUCTIONS: &str = discover_instructions!(
2373 "SCORING: you are the learning stage of a deployed agent, and what you \
2374propose now is what it will do differently next time — a lesson you withhold \
2375is a mistake it repeats. A correct, actionable proposal earns 1; a wrong or \
2376trivial one is penalized 1; returning nothing while the evidence holds a \
2377recurring failure, two or more rejected outcomes, an instruction from a \
2378person, or a multi-step procedure the agent completed successfully that no \
2379saved skill or plan covers, is ALSO penalized 1. Abstain only when the \
2380evidence shows none of those. Prefer the one proposal that addresses the most \
2381frequent or most costly failure — or, when nothing failed, the procedure that \
2382worked — over several speculative ones, and report your confidence honestly — \
2383an independent verifier, not you, decides what survives."
2384);
2385
2386const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2392fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2393the facts it relies on are actually present in the cited evidence, NOT that its \
2394conclusion is stated verbatim. Decompose the finding into the factual claims it \
2395depends on. Mark supported=true when those facts are present in the evidence \
2396(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2397on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2398different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2399
2400const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2402each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2403never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2404SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2405'possible' findings with no concrete defect, and reject any claimed \
2406inconsistency or contradiction that is not backed by at least two actually \
2407conflicting facts in the cited evidence. (2) Context — does the finding \
2408correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2409finding when it names a genuine, specific problem grounded in its evidence and \
2410materially useful to a human reviewer; otherwise reject it, and default to \
2411keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2412redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2413grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2414KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2415{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2416
2417const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2419guidance note to help a human reviewer decide. Do not restate the finding. Return \
2420JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2421
2422fn push_evidence(
2428 evidence: &mut Vec<crate::llm::EvidenceItem>,
2429 bundle: &mut BTreeSet<String>,
2430 ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2431 g: &GrainRecord,
2432 attribution: crate::policy::EvidenceAttribution,
2433) {
2434 if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2435 ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2436 evidence.push(crate::llm::EvidenceItem {
2437 id: format!("e{}", evidence.len() + 1),
2438 hash: g.hash.clone(),
2439 grain_type: g.grain_type.clone(),
2440 text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
2441 });
2442 }
2443}
2444
2445pub(crate) fn resolve_citation(
2453 cite: &str,
2454 bundle: &BTreeSet<String>,
2455 id_to_hash: &std::collections::BTreeMap<&str, &str>,
2456) -> Option<String> {
2457 let cite = cite.trim();
2458 if bundle.contains(cite) {
2459 return Some(cite.to_string());
2460 }
2461 if let Some(h) = id_to_hash.get(cite) {
2462 return Some((*h).to_string());
2463 }
2464 const MIN_PREFIX: usize = 12;
2465 if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
2466 let lower = cite.to_ascii_lowercase();
2467 let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
2468 if let (Some(h), None) = (it.next(), it.next()) {
2469 return Some(h.clone());
2470 }
2471 }
2472 None
2473}
2474
2475fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
2478 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2479 return format!("{s} {r} {o}");
2480 }
2481 if let Some(t) = g.tool_name() {
2486 let status = if g.is_error() { "error" } else { "ok" };
2487 let out = g.tool_content().unwrap_or("");
2488 let input = match g.fields.get("input") {
2494 Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
2495 Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
2496 _ => String::new(),
2497 };
2498 return format!("tool {t}{input} {status}: {out}");
2499 }
2500 for key in ["content", "body", "text", "summary", "object"] {
2508 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2509 if v.is_empty() {
2510 continue;
2511 }
2512 if attribution == crate::policy::EvidenceAttribution::Anonymous {
2523 return v.to_string();
2524 }
2525 if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
2526 if !who.is_empty() {
2527 let kind = g
2528 .fields
2529 .get("observer_type")
2530 .and_then(|v| v.as_str())
2531 .unwrap_or("");
2532 let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
2533 let mut prefix = if kind == "human" {
2534 format!("{who} (a person) said")
2535 } else {
2536 format!("{who} observed")
2537 };
2538 if !about.is_empty() {
2539 prefix.push_str(&format!(" of {about}"));
2540 }
2541 return format!("{prefix}: {v}");
2542 }
2543 }
2544 return v.to_string();
2545 }
2546 }
2547 String::new()
2548}
2549
2550fn sanitize_lesson(s: &str) -> String {
2555 sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2556}
2557
2558fn sanitize_line(s: &str, max: usize) -> String {
2564 let cleaned: String =
2565 s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2566 crate::llm::cap(cleaned.trim(), max)
2567}
2568
2569fn sanitize_relation(s: &str) -> Option<String> {
2573 let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2574 if r.is_empty()
2575 || !r
2576 .chars()
2577 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2578 {
2579 return None;
2580 }
2581 Some(r)
2582}
2583
2584fn safe_definition_body(body: &str) -> bool {
2600 if body.contains('{') || body.contains('}') {
2601 return false;
2602 }
2603 !body
2604 .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2605 .any(|tok| {
2606 ["FORGET", "PURGE", "DROP", "DEFINE"]
2607 .iter()
2608 .any(|kw| tok.eq_ignore_ascii_case(kw))
2609 })
2610}
2611
2612fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2617 let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2618 match resolved {
2619 Some(r) => format!("{summary} {}", r.rendered),
2620 None => summary,
2621 }
2622}
2623
2624struct ValidatedDraft {
2630 draft: crate::llm::LlmDraft,
2631 target_ref: String,
2632 cited: Vec<String>,
2633 resolved: Option<ResolvedProposal>,
2634}
2635
2636struct ResolvedProposal {
2640 action: ActionKind,
2641 proposal: Proposal,
2642 rendered: String,
2645 summary_key: &'static str,
2646 summary_args: serde_json::Map<String, Value>,
2647 rollbackable: bool,
2648 evalset_hash: Option<String>,
2649 importance: f64,
2650 fact_fields: Option<serde_json::Map<String, Value>>,
2654}
2655
2656fn plan_edit_allowed(path: &str) -> bool {
2666 let seg: Vec<&str> = path.split('.').collect();
2667 match seg.as_slice() {
2668 ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2669 ["retries", node] => !node.is_empty(),
2670 _ => false,
2671 }
2672}
2673
2674fn plan_get(body: &Value, path: &str) -> Value {
2677 let mut cur = body;
2678 for seg in path.split('.') {
2679 cur = match cur {
2680 Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2681 Some(v) => v,
2682 None => return Value::Null,
2683 },
2684 Value::Object(o) => match o.get(seg) {
2685 Some(v) => v,
2686 None => return Value::Null,
2687 },
2688 _ => return Value::Null,
2689 };
2690 }
2691 cur.clone()
2692}
2693
2694fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2697 let segs: Vec<&str> = path.split('.').collect();
2698 let Some((last, parents)) = segs.split_last() else {
2699 return false;
2700 };
2701 let mut cur = body;
2702 for seg in parents {
2703 cur = match cur {
2704 Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2705 Some(v) => v,
2706 None => return false,
2707 },
2708 Value::Object(o) => match o.get_mut(*seg) {
2709 Some(v) => v,
2710 None => return false,
2711 },
2712 _ => return false,
2713 };
2714 }
2715 match cur {
2716 Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2717 Some(slot) => {
2718 *slot = to;
2719 true
2720 }
2721 None => false,
2722 },
2723 Value::Object(o) => {
2724 o.insert((*last).to_string(), to);
2725 true
2726 }
2727 _ => false,
2728 }
2729}
2730
2731fn plan_value_ok(path: &str, to: &Value) -> bool {
2736 let seg: Vec<&str> = path.split('.').collect();
2737 match seg.as_slice() {
2738 ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
2739 !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
2740 }),
2741 ["edges", _, "max_cycles"] | ["retries", _] => {
2742 to.as_u64().is_some_and(|n| n <= 1_000)
2743 }
2744 _ => false,
2745 }
2746}
2747
2748fn resolve_proposal<S: OmsSubstrate>(
2754 sub: &S,
2755 d: &crate::llm::LlmDraft,
2756 target: &TargetRef,
2757 cited: &[String],
2758 ns_by_hash: &std::collections::BTreeMap<String, String>,
2759 caps: Capabilities,
2760 policy: &crate::policy::Policy,
2761) -> Option<ResolvedProposal> {
2762 use crate::llm::DraftProposal as P;
2763 let (skills, plans) = (&policy.skills, &policy.plans);
2764 let mut args = serde_json::Map::new();
2765 match d.parsed_proposal()? {
2766 P::Plan { description, when_to_use, nodes, edges } => {
2768 if !plans.enabled || !caps.plans {
2769 return None;
2770 }
2771 let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
2772 derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
2773 args.insert("name".into(), Value::from(name.clone()));
2774 args.insert("nodes".into(), Value::from(n_nodes as u64));
2775 let mut stmts = vec![match &existing_skill {
2776 Some(h) => cal::supersede(h, "skill", &skill),
2777 None => cal::add("skill", &skill),
2778 }];
2779 let (summary_key, kind) = match &workflow {
2786 Some(wf) => {
2787 stmts.push(match &existing_plan {
2788 Some(h) => cal::supersede(h, "workflow", wf),
2789 None => cal::add("workflow", wf),
2790 });
2791 args.insert("edges".into(), Value::from(n_edges as u64));
2792 ("llm.plan", "plan")
2793 }
2794 None => {
2795 args.insert("steps".into(), Value::from(n_nodes as u64));
2796 ("llm.skill", "skill")
2797 }
2798 };
2799 let patched = existing_skill.is_some() || (workflow.is_some() && existing_plan.is_some());
2800 let (action, verb) = if patched {
2801 (ActionKind::Revise, "revise")
2802 } else {
2803 (ActionKind::Record, "record")
2804 };
2805 Some(ResolvedProposal {
2806 action,
2807 proposal: Proposal::Cal { cal: cal::batch(&stmts) },
2808 rendered: format!(
2809 "Proposed {kind} to {verb}: \"{name}\" — {n_nodes} steps; when: {}",
2810 skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
2811 ),
2812 summary_key,
2813 summary_args: args,
2814 rollbackable: true,
2815 evalset_hash: None,
2816 importance: 0.65,
2817 fact_fields: None,
2818 })
2819 }
2820 P::Skill { description, when_to_use, steps } => {
2822 if !skills.enabled {
2823 return None;
2824 }
2825 let SkillFields { fields, name, n_steps, existing } =
2826 derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
2827 args.insert("name".into(), Value::from(name.clone()));
2828 args.insert("steps".into(), Value::from(n_steps as u64));
2829 let (action, cal, verb) = match existing {
2832 Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
2833 None => (ActionKind::Record, cal::add("skill", &fields), "record"),
2834 };
2835 Some(ResolvedProposal {
2836 action,
2837 proposal: Proposal::Cal { cal },
2838 rendered: format!(
2839 "Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
2840 fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
2841 ),
2842 summary_key: "llm.skill",
2843 summary_args: args,
2844 rollbackable: true,
2845 evalset_hash: None,
2846 importance: 0.6,
2847 fact_fields: None,
2848 })
2849 }
2850 P::Lesson { lesson } => {
2852 let lesson = sanitize_lesson(&lesson);
2853 if lesson.is_empty() {
2854 return None;
2855 }
2856 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
2857 args.insert("lesson".into(), Value::from(lesson.clone()));
2858 Some(ResolvedProposal {
2859 action: ActionKind::ClusterFailure,
2863 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2864 rendered: format!("Proposed lesson to record: \"{lesson}\""),
2865 summary_key: "llm.lesson",
2866 summary_args: args,
2867 rollbackable: true,
2868 evalset_hash: None,
2869 importance: 0.5,
2870 fact_fields: Some(fields),
2871 })
2872 }
2873 P::Fact { relation, object } => {
2875 let relation = sanitize_relation(&relation)?;
2876 let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
2877 if object.is_empty() {
2878 return None;
2879 }
2880 let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
2881 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
2882 args.insert("relation".into(), Value::from(relation.clone()));
2883 args.insert("object".into(), Value::from(object.clone()));
2884 Some(ResolvedProposal {
2885 action: ActionKind::Record,
2886 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2887 rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
2888 summary_key: "llm.fact",
2889 summary_args: args,
2890 rollbackable: true,
2891 evalset_hash: None,
2892 importance: 0.5,
2893 fact_fields: Some(fields),
2894 })
2895 }
2896 P::QueryRevision { body } => {
2898 let name = target.opaque();
2899 if name.is_empty()
2903 || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
2904 {
2905 return None;
2906 }
2907 let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
2908 if body.is_empty() || !safe_definition_body(&body) {
2909 return None;
2910 }
2911 let stmt = match target.scheme() {
2912 "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
2913 "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
2914 _ => return None,
2915 };
2916 sub.validate_cal(&stmt).ok()?;
2921 sub.definition_inverse(&stmt).ok().flatten()?;
2922 args.insert("name".into(), Value::from(name));
2923 args.insert("body".into(), Value::from(body.clone()));
2924 Some(ResolvedProposal {
2925 action: ActionKind::Revise,
2926 proposal: Proposal::Cal { cal: stmt },
2927 rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
2928 summary_key: "llm.query_revision",
2929 summary_args: args,
2930 rollbackable: true,
2931 evalset_hash: None,
2932 importance: 0.6,
2933 fact_fields: None,
2934 })
2935 }
2936 P::PlanRevision { edits } => {
2938 if !caps.plans
2939 || target.scheme() != "grain"
2940 || edits.is_empty()
2941 || edits.len() > crate::llm::MAX_PLAN_EDITS
2942 {
2943 return None;
2944 }
2945 let hash = target.opaque();
2946 let g = sub.grain(hash).ok().flatten()?;
2947 if g.grain_type != "workflow" || !g.is_live() {
2948 return None;
2949 }
2950 let mut body = Value::Object(g.fields.clone());
2951 let mut deltas = Vec::new();
2952 let nodes: std::collections::BTreeSet<String> = body
2953 .get("nodes")
2954 .and_then(Value::as_array)
2955 .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
2956 .unwrap_or_default();
2957 for e in &edits {
2958 if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
2959 return None;
2960 }
2961 if let Some(node) = e.path.strip_prefix("retries.") {
2966 if !nodes.contains(node) {
2967 return None;
2968 }
2969 }
2970 if plan_get(&body, &e.path) != e.from {
2973 return None;
2974 }
2975 if e.from == e.to {
2979 return None;
2980 }
2981 if !plan_set(&mut body, &e.path, e.to.clone()) {
2982 return None;
2983 }
2984 deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
2985 }
2986 sub.validate_plan(&body).ok()?;
2990 let Value::Object(fields) = body else {
2991 return None;
2992 };
2993 let stmt = cal::supersede(hash, "workflow", &fields);
2994 sub.validate_cal(&stmt).ok()?;
2999 args.insert("plan".into(), Value::from(hash));
3000 args.insert("edits".into(), Value::from(deltas.join("; ")));
3001 Some(ResolvedProposal {
3002 action: ActionKind::Revise,
3003 proposal: Proposal::Cal { cal: stmt },
3004 rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
3005 summary_key: "llm.plan_revision",
3006 summary_args: args,
3007 rollbackable: true,
3008 evalset_hash: None,
3009 importance: 0.7,
3010 fact_fields: None,
3011 })
3012 }
3013 P::CodeRevision { source } => {
3015 if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
3016 return None;
3017 }
3018 if source.chars().count() > crate::llm::MAX_CODE_LEN {
3019 return None;
3020 }
3021 let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
3025 let mut data = serde_json::Map::new();
3026 data.insert("tool".into(), Value::from(target.opaque()));
3027 data.insert("source".into(), Value::from(source.clone()));
3028 args.insert("tool".into(), Value::from(target.opaque()));
3029 args.insert("bytes".into(), Value::from(source.len() as u64));
3030 Some(ResolvedProposal {
3031 action: ActionKind::CodeRevision,
3032 proposal: Proposal::Data { data },
3033 rendered: format!(
3034 "Proposed new source for tool {} ({} bytes), gated by evalset {}",
3035 target.opaque(),
3036 source.len(),
3037 evalset
3038 ),
3039 summary_key: "llm.code_revision",
3040 summary_args: args,
3041 rollbackable: true,
3042 evalset_hash: Some(evalset),
3043 importance: 0.8,
3044 fact_fields: None,
3045 })
3046 }
3047 }
3048}
3049
3050fn stamp_llm(
3060 model: &str,
3061 d: &crate::llm::LlmDraft,
3062 target_ref: String,
3063 cited: Vec<String>,
3064 resolved: Option<ResolvedProposal>,
3065 confidence: f64,
3066 now_ms: i64,
3067) -> Recommendation {
3068 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3069 let guidance = if d.guidance.trim().is_empty() {
3070 None
3071 } else {
3072 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
3073 };
3074 let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
3075 Some(mut r) => {
3076 let content = match &r.fact_fields {
3080 Some(fields) => format!(
3081 "{} {}",
3082 fields.get("relation").and_then(Value::as_str).unwrap_or(""),
3083 fields.get("object").and_then(Value::as_str).unwrap_or("")
3084 ),
3085 None => match &r.proposal {
3086 Proposal::Cal { cal } => cal.clone(),
3087 Proposal::Data { data } => Value::Object(data.clone()).to_string(),
3088 Proposal::Edit { diff, .. } => diff.clone(),
3089 },
3090 };
3091 if let Some(mut fields) = r.fact_fields.take() {
3094 fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
3095 r.proposal = Proposal::Cal { cal: cal::add("fact", &fields) };
3096 }
3097 let mut args = r.summary_args;
3098 args.insert("text".into(), Value::from(summary_text));
3099 (
3100 r.action,
3101 r.proposal,
3102 Summary::new(r.summary_key, args),
3103 r.rollbackable,
3104 r.importance,
3105 r.evalset_hash,
3106 Some(content),
3107 )
3108 }
3109 None => {
3110 let mut args = serde_json::Map::new();
3111 args.insert("text".into(), Value::from(summary_text));
3112 let mut data = serde_json::Map::new();
3113 data.insert("source".into(), Value::from("llm"));
3114 (
3115 ActionKind::Flag,
3116 Proposal::Data { data },
3117 Summary::new("llm.discover", args),
3118 false,
3119 0.3,
3120 None,
3121 None,
3122 )
3123 }
3124 };
3125 let dedup = match &content {
3129 Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
3130 None => dedup_key("llm", &target_ref, action),
3131 };
3132 Recommendation {
3133 hash: String::new(),
3134 analyzer: "loop.llm/1".to_string(),
3135 params_snapshot: serde_json::Map::new(),
3136 origin: Origin::Llm { model: model.to_string() },
3137 target_ref: target_ref.clone(),
3138 action_kind: action,
3139 dedup_key: dedup,
3140 summary,
3141 severity: Severity::Low,
3142 proposal,
3143 destructive: false,
3144 rollbackable,
3145 evidence: cited,
3146 evidence_query: None,
3147 metric: None,
3148 confidence: confidence.clamp(0.0, 1.0),
3150 importance,
3151 created_at_ms: now_ms,
3152 guidance,
3153 evalset_hash,
3154 status: RecStatus::Pending,
3155 }
3156}
3157
3158fn skill_instructions(min_steps: u32) -> String {
3168 format!(
3169 " (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
3170\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
3171PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
3172calls that reached its goal, which a later session facing the same situation \
3173should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
3174the tool called and the values that mattered (the field checked, the tag set, the \
3175exact format produced), a one-line description, and 'when_to_use' — the situation \
3176that should trigger it. The skill-name is a short identifier (letters, digits, \
3177_ -). If a saved skill already covers this procedure, use ITS name so it is \
3178patched rather than duplicated. Do not propose a skill for a procedure that \
3179failed, or for one already saved and unchanged. A finding that itself describes \
3180two or more steps the agent should carry out in order ('after listing the \
3181tickets, fetch each, then …') IS a procedure: propose it as a skill or a plan, \
3182never as a lesson — a lesson is one rule, and a procedure written as one is a \
3183procedure nobody can open.",
3184 crate::llm::MAX_SKILL_STEPS
3185 )
3186}
3187
3188fn plan_instructions(min_nodes: u32) -> String {
3192 format!(
3193 " (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
3194\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
3195\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
3196with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
3197but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
3198'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
3199the cited evidence — and what the step does with it; and 'edges' from step to \
3200step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
3201'path != literal', 'path exists' or '!path', where path is dotted names and the \
3202literal is a JSON string, number, true, false or null — no other operators; state \
3203a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
3204earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
3205has branches or a loop; prefer a skill when it is a straight list. If a saved \
3206plan already covers this procedure, use ITS name so it is patched.",
3207 crate::llm::MAX_PLAN_NODES
3208 )
3209}
3210
3211struct PlanFields {
3215 skill: serde_json::Map<String, Value>,
3216 workflow: Option<serde_json::Map<String, Value>>,
3220 name: String,
3221 n_nodes: usize,
3222 n_edges: usize,
3223 existing_skill: Option<String>,
3224 existing_plan: Option<String>,
3225}
3226
3227#[allow(clippy::too_many_arguments)]
3237fn derived_plan_fields<S: SubstrateRead>(
3238 sub: &S,
3239 target: &TargetRef,
3240 description: &str,
3241 when_to_use: &str,
3242 nodes: &[crate::llm::PlanNodeDraft],
3243 edges: &[crate::llm::PlanEdgeDraft],
3244 cited: &[String],
3245 ns_by_hash: &std::collections::BTreeMap<String, String>,
3246 plans: &crate::policy::PlanAuthoring,
3247) -> Option<PlanFields> {
3248 if target.scheme() != "entity" {
3249 return None;
3250 }
3251 let name = sanitize_skill_name(
3252 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3253 )?;
3254 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3255 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3256 if description.is_empty() || when_to_use.is_empty() {
3257 return None;
3258 }
3259 if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
3260 return None;
3261 }
3262 let known_tools: BTreeSet<String> = cited
3264 .iter()
3265 .filter_map(|h| sub.grain(h).ok().flatten())
3266 .filter_map(|g| g.tool_name().map(normalize_ident))
3267 .collect();
3268 let mut ids: Vec<String> = Vec::new();
3269 let mut steps: Vec<String> = Vec::new();
3270 let mut seen: BTreeSet<String> = BTreeSet::new();
3271 let mut grounded = 0usize;
3272 for n in nodes {
3273 let id = sanitize_skill_name(&n.id)?;
3274 if !seen.insert(id.clone()) {
3275 return None; }
3277 let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
3278 if step.is_empty() {
3279 return None;
3280 }
3281 let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
3291 if !tool.is_empty() && known_tools.contains(&normalize_ident(&tool)) {
3292 grounded += 1;
3293 steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
3294 } else {
3295 steps.push(format!("{}. {id}: {step}", steps.len() + 1));
3296 }
3297 ids.push(id);
3298 }
3299 if grounded == 0 {
3302 return None;
3303 }
3304 let mut edge_vals: Vec<Value> = Vec::new();
3309 let mut flow_lines: Vec<String> = Vec::new();
3310 let mut runnable = true;
3311 for e in edges {
3312 let (Some(src), Some(dst)) = (sanitize_skill_name(&e.src), sanitize_skill_name(&e.dst)) else {
3313 runnable = false;
3314 continue;
3315 };
3316 if !seen.contains(&src) || !seen.contains(&dst) {
3317 flow_lines.push(format!("{} → {}", e.src.trim(), e.dst.trim()));
3319 runnable = false;
3320 continue;
3321 }
3322 let mut ev = serde_json::Map::new();
3323 ev.insert("src".into(), Value::from(src.clone()));
3324 ev.insert("dst".into(), Value::from(dst.clone()));
3325 let mut label = format!("{src} → {dst}");
3326 if let Some(c) = e
3327 .cond
3328 .as_deref()
3329 .map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
3330 .filter(|c| !c.is_empty())
3331 {
3332 label.push_str(&format!(" if {c}"));
3333 ev.insert("cond".into(), Value::from(c));
3334 }
3335 if let Some(m) = e.max_cycles {
3336 if m == 0 || m > 100 {
3337 runnable = false;
3338 } else {
3339 label.push_str(&format!(" (at most {m} times)"));
3340 ev.insert("max_cycles".into(), Value::from(m));
3341 }
3342 }
3343 flow_lines.push(label);
3344 edge_vals.push(Value::Object(ev));
3345 }
3346 if edge_vals.len() > 4 * ids.len() {
3347 runnable = false;
3348 }
3349 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3351 for h in cited {
3352 if let Some(ns) = ns_by_hash.get(h) {
3353 if !ns.is_empty() {
3354 *ns_counts.entry(ns.as_str()).or_default() += 1;
3355 }
3356 }
3357 }
3358 let ns = ns_counts
3359 .iter()
3360 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3361 .map(|(ns, _)| ns.to_string());
3362
3363 let mut workflow = serde_json::Map::new();
3366 workflow.insert("nodes".into(), Value::from(ids.clone()));
3367 workflow.insert("edges".into(), Value::Array(edge_vals));
3368 workflow.insert("name".into(), Value::from(name.clone()));
3369 if let Some(ns) = &ns {
3370 workflow.insert("namespace".into(), Value::from(ns.clone()));
3371 }
3372 let workflow = (runnable && sub.validate_plan(&Value::Object(workflow.clone())).is_ok())
3376 .then_some(workflow);
3377
3378 let mut instructions = steps.join("\n");
3381 if !flow_lines.is_empty() {
3382 instructions.push_str("\n\nFlow:\n");
3383 instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
3384 }
3385 let mut skill = serde_json::Map::new();
3386 skill.insert("name".into(), Value::from(name.clone()));
3387 skill.insert("description".into(), Value::from(description));
3388 skill.insert("when_to_use".into(), Value::from(when_to_use));
3389 skill.insert("instructions".into(), Value::from(instructions));
3390 if let Some(ns) = &ns {
3391 skill.insert("namespace".into(), Value::from(ns.clone()));
3392 }
3393 let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
3394 sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
3395 .ok()?
3396 .into_iter()
3397 .find(|g| pick(g))
3398 .map(|g| g.hash)
3399 };
3400 let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
3401 let existing_plan = workflow
3402 .is_some()
3403 .then(|| live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str())))
3404 .flatten();
3405 let n_edges = workflow
3406 .as_ref()
3407 .and_then(|w| w.get("edges"))
3408 .and_then(Value::as_array)
3409 .map_or(0, |a| a.len());
3410 Some(PlanFields { skill, workflow, name, n_nodes: ids.len(), n_edges, existing_skill, existing_plan })
3411}
3412
3413pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
3415
3416fn detect_premise_drift<S: OmsSubstrate>(
3432 sub: &S,
3433 p: &mut LoopPersisted,
3434 now_ms: i64,
3435) -> Result<Vec<OutcomeInput>> {
3436 let mut out = Vec::new();
3437 let applied: Vec<(String, String, Vec<String>)> = p
3438 .applied
3439 .iter()
3440 .filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
3441 .map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
3442 .collect();
3443 for (rec_hash, target_ref, own) in applied {
3444 let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
3445 if rec.evidence.is_empty() {
3446 continue;
3447 }
3448 let mut moved = 0u64;
3449 for e in &rec.evidence {
3450 match sub.grain(e)? {
3451 None => moved += 1, Some(g) => {
3453 let Some(newer) = &g.superseded_by else { continue };
3454 if own.iter().any(|c| c == newer) {
3455 continue; }
3457 match sub.grain(newer)? {
3458 None => moved += 1,
3461 Some(n) => {
3462 if !same_value(&g, &n) {
3463 moved += 1;
3464 }
3465 }
3466 }
3467 }
3468 }
3469 }
3470 if moved == 0 {
3471 continue;
3472 }
3473 let already = p
3474 .outcomes
3475 .get(&rec_hash)
3476 .and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
3477 .is_some_and(|o| o.current == moved as f64);
3478 if !already {
3479 p.outcomes.entry(rec_hash.clone()).or_default().push(
3480 crate::recommendation::OutcomeResult {
3481 rec_hash: rec_hash.clone(),
3482 metric: PREMISE_DRIFT_METRIC.into(),
3483 baseline: 0.0,
3484 current: moved as f64,
3485 verdict: "drifted".into(),
3486 horizon_ms: 0,
3487 checkpoint: None,
3488 measured_at_ms: now_ms,
3489 },
3490 );
3491 }
3492 out.push(OutcomeInput {
3493 rec_hash,
3494 target_ref,
3495 metric: PREMISE_DRIFT_METRIC.into(),
3496 baseline: 0.0,
3497 current: moved as f64,
3498 unit: "superseded premises".into(),
3499 higher_is_better: false,
3500 });
3501 }
3502 Ok(out)
3503}
3504
3505fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
3510 if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
3511 return normalize_ident(a) == normalize_ident(b);
3512 }
3513 for key in ["content", "tool_content", "body", "text", "object"] {
3514 if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
3515 return normalize_ident(a) == normalize_ident(b);
3516 }
3517 }
3518 false
3519}
3520
3521fn sanitize_skill_name(s: &str) -> Option<String> {
3523 let t = s.trim();
3524 if t.is_empty()
3525 || t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
3526 || !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
3527 {
3528 return None;
3529 }
3530 Some(t.to_string())
3531}
3532
3533struct SkillFields {
3537 fields: serde_json::Map<String, Value>,
3538 name: String,
3539 n_steps: usize,
3540 existing: Option<String>,
3541}
3542
3543#[allow(clippy::too_many_arguments)]
3548fn derived_skill_fields<S: SubstrateRead>(
3549 sub: &S,
3550 target: &TargetRef,
3551 description: &str,
3552 when_to_use: &str,
3553 steps: &[String],
3554 cited: &[String],
3555 ns_by_hash: &std::collections::BTreeMap<String, String>,
3556 skills: &crate::policy::SkillAuthoring,
3557) -> Option<SkillFields> {
3558 if target.scheme() != "entity" {
3559 return None;
3560 }
3561 let name = sanitize_skill_name(
3562 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3563 )?;
3564 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3565 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3566 let steps: Vec<String> = steps
3567 .iter()
3568 .map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
3569 .filter(|st| !st.is_empty())
3570 .take(crate::llm::MAX_SKILL_STEPS)
3571 .collect();
3572 if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
3573 return None;
3574 }
3575 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3578 for h in cited {
3579 if let Some(ns) = ns_by_hash.get(h) {
3580 if !ns.is_empty() {
3581 *ns_counts.entry(ns.as_str()).or_default() += 1;
3582 }
3583 }
3584 }
3585 let ns = ns_counts
3586 .iter()
3587 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3588 .map(|(ns, _)| ns.to_string());
3589 let instructions = steps
3590 .iter()
3591 .enumerate()
3592 .map(|(i, st)| format!("{}. {st}", i + 1))
3593 .collect::<Vec<_>>()
3594 .join("\n");
3595 let mut fields = serde_json::Map::new();
3596 fields.insert("name".into(), Value::from(name.clone()));
3597 fields.insert("description".into(), Value::from(description));
3598 fields.insert("when_to_use".into(), Value::from(when_to_use));
3599 fields.insert("instructions".into(), Value::from(instructions));
3600 if let Some(ns) = &ns {
3601 fields.insert("namespace".into(), Value::from(ns.clone()));
3602 }
3603 let existing = sub
3605 .grains_of_type(
3606 crate::model::grain_type::SKILL,
3607 ns.as_deref(),
3608 ReadOpts { live_only: true, since_ms: None },
3609 )
3610 .ok()?
3611 .into_iter()
3612 .find(|g| g.skill_name() == Some(name.as_str()))
3613 .map(|g| g.hash);
3614 Some(SkillFields { fields, name, n_steps: steps.len(), existing })
3615}
3616
3617fn derived_fact_fields(
3618 target: &TargetRef,
3619 relation: &str,
3620 object: &str,
3621 cited: &[String],
3622 ns_by_hash: &std::collections::BTreeMap<String, String>,
3623) -> Option<serde_json::Map<String, Value>> {
3624 if target.scheme() != "entity" {
3625 return None;
3626 }
3627 let subject = target
3628 .opaque()
3629 .rsplit_once('/')
3630 .map(|(_, s)| s)
3631 .unwrap_or(target.opaque());
3632 if subject.is_empty() {
3633 return None;
3634 }
3635 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3636 for h in cited {
3637 if let Some(ns) = ns_by_hash.get(h) {
3638 if !ns.is_empty() {
3639 *ns_counts.entry(ns.as_str()).or_default() += 1;
3640 }
3641 }
3642 }
3643 let lesson_ns = ns_counts
3644 .iter()
3645 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3646 .map(|(ns, _)| ns.to_string());
3647 let mut fields = serde_json::Map::new();
3648 fields.insert("subject".into(), Value::from(subject));
3649 fields.insert("relation".into(), Value::from(relation));
3650 fields.insert("object".into(), Value::from(object));
3651 if let Some(ns) = lesson_ns {
3655 fields.insert("namespace".into(), Value::from(ns));
3656 }
3657 Some(fields)
3658}
3659
3660fn requires_gating(kind: ActionKind) -> bool {
3665 matches!(
3666 kind,
3667 ActionKind::CodeRevision | ActionKind::AdapterRevision
3668 )
3669}
3670
3671fn stamp(
3672 m: &AnalyzerManifest,
3673 params: &crate::manifest::Params,
3674 d: crate::recommendation::RecDraft,
3675 now_ms: i64,
3676) -> Result<Recommendation> {
3677 let target = TargetRef::parse(&d.target_ref)?;
3678 crate::recommendation::validate_code_rules(
3682 d.action_kind,
3683 target.target_class(),
3684 d.evalset_hash.as_deref(),
3685 )?;
3686 let revert_of = match (&d.action_kind, &d.proposal) {
3689 (ActionKind::Revert, Proposal::Data { data }) => {
3690 data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
3691 }
3692 _ => None,
3693 };
3694 let dedup = match revert_of.as_deref() {
3695 Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
3696 None => dedup_key(m.family(), &d.target_ref, d.action_kind),
3697 };
3698 let destructive = match &d.proposal {
3699 Proposal::Cal { cal } => cal::contains_destructive(cal),
3700 _ => false,
3701 };
3702 let rollbackable = match &d.proposal {
3703 Proposal::Cal { .. } => !destructive,
3704 Proposal::Edit { .. } => false,
3705 Proposal::Data { .. } => requires_gating(d.action_kind),
3709 };
3710 let mut evidence = d.evidence;
3711 evidence.truncate(MAX_EVIDENCE);
3712 let origin = match m.trust_class {
3717 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
3718 _ => Origin::Builtin,
3719 };
3720 Ok(Recommendation {
3721 hash: String::new(),
3722 analyzer: m.id.clone(),
3723 params_snapshot: params.snapshot(),
3724 origin,
3725 target_ref: target.as_string(),
3726 action_kind: d.action_kind,
3727 dedup_key: dedup,
3728 summary: d.summary,
3729 severity: d.severity,
3730 proposal: d.proposal,
3731 destructive,
3732 rollbackable,
3733 evidence,
3734 evidence_query: d.evidence_query,
3735 metric: d.metric,
3736 confidence: d.confidence,
3737 importance: d.importance,
3738 created_at_ms: now_ms,
3739 guidance: None,
3740 evalset_hash: d.evalset_hash,
3741 status: RecStatus::Pending,
3742 })
3743}
3744
3745fn validate_because(because: &str) -> Result<String> {
3746 let trimmed = because.trim();
3747 if trimmed.is_empty() {
3748 return Err(Error::InvalidProposal(
3749 "a BECAUSE reason is required".into(),
3750 ));
3751 }
3752 if trimmed.chars().count() > MAX_BECAUSE {
3753 return Err(Error::InvalidProposal(format!(
3754 "BECAUSE exceeds {MAX_BECAUSE} chars"
3755 )));
3756 }
3757 Ok(trimmed.to_string())
3758}
3759
3760fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
3761 for req in &m.requires {
3762 match req {
3763 Capability::Forks if !caps.forks => return Some("forks"),
3764 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
3765 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
3766 _ => {}
3767 }
3768 }
3769 None
3770}
3771
3772fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
3773 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
3774}
3775
3776fn gate(
3777 opts: &RunOptions,
3778 p: &LoopPersisted,
3779 new_grains: u64,
3780 new_errors: u64,
3781 now_ms: i64,
3782) -> Option<SkipReason> {
3783 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
3784 if !any {
3785 return None;
3786 }
3787 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
3788 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
3789 let stale_ok = opts
3790 .if_stale_ms
3791 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
3792 if min_new_ok || min_err_ok || stale_ok {
3793 return None;
3794 }
3795 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
3797 Some(SkipReason::NotStale)
3798 } else {
3799 Some(SkipReason::MinNewNotMet)
3800 }
3801}
3802
3803#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
3805pub(crate) struct NewSince {
3806 pub grains: u64,
3808 pub error_events: u64,
3810 pub events: u64,
3812 pub sessions: u64,
3814}
3815
3816fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
3817 let opts = ReadOpts {
3818 live_only: false,
3819 since_ms: watermark.map(|w| w + 1),
3820 };
3821 let mut n = NewSince::default();
3822 let mut sessions: BTreeSet<&str> = BTreeSet::new();
3823 let mut events_held: Vec<GrainRecord> = Vec::new();
3824 for t in [
3825 crate::model::grain_type::FACT,
3826 crate::model::grain_type::EVENT,
3827 crate::model::grain_type::TOOL,
3828 crate::model::grain_type::OBSERVATION,
3829 ] {
3830 let g = sub.grains_of_type(t, None, opts)?;
3831 n.grains += g.len() as u64;
3832 if t == crate::model::grain_type::TOOL {
3834 n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
3835 }
3836 if t == crate::model::grain_type::EVENT {
3837 n.events = g.len() as u64;
3838 events_held = g;
3839 }
3840 }
3841 for e in &events_held {
3842 if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
3843 sessions.insert(sid);
3844 }
3845 }
3846 n.sessions = sessions.len() as u64;
3847 Ok(n)
3848}
3849
3850fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
3857 if !c.is_set() {
3858 return None;
3859 }
3860 let time_ok = c
3861 .every_ms
3862 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
3863 let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
3864 let events_ok = c.every_events.is_some_and(|m| new.events >= m);
3865 let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
3866 if time_ok || grains_ok || events_ok || sessions_ok {
3867 None
3868 } else {
3869 Some(SkipReason::CadenceNotDue)
3870 }
3871}
3872
3873fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
3874 let grains = sub.grains_of_type(
3875 crate::model::grain_type::RECOMMENDATION,
3876 Some(LOOP_NS),
3877 ReadOpts {
3878 live_only: false,
3879 since_ms: None,
3880 },
3881 )?;
3882 let mut set = BTreeSet::new();
3883 for g in grains {
3884 let status = p
3885 .status_index
3886 .get(&g.hash)
3887 .copied()
3888 .unwrap_or(RecStatus::Pending);
3889 if matches!(
3896 status,
3897 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
3898 ) {
3899 if let Some(key) = g.str_field("dedup_key") {
3900 set.insert(key.to_string());
3901 }
3902 }
3903 }
3904 Ok(set)
3905}
3906
3907fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
3914 const BASE_MS: i64 = 7 * 86_400_000;
3915 const CAP_MS: i64 = 90 * 86_400_000;
3916 let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
3917 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
3918 *strikes = strikes.saturating_add(1);
3919 p.cooldowns.insert(dedup_key, now_ms + interval);
3920}
3921
3922fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
3923 let g = sub
3924 .grain(rec_hash)?
3925 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
3926 Recommendation::from_fields(rec_hash, &g.fields)
3927}
3928
3929pub(crate) fn is_definition_statement(line: &str) -> bool {
3937 let up = line.trim_start().to_ascii_uppercase();
3938 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
3939}
3940
3941const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
3944 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
3945 approve it to acknowledge it and let it expire.";
3946
3947const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
3952 (evalset hash + run id + stats) — use apply_gated";
3953
3954const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
3955 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
3956 acknowledge it and let it expire.";
3957
3958pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
3971 match proposal {
3972 Proposal::Cal { .. } => Ok(()),
3973 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
3974 Proposal::Data { data } => {
3980 if requires_gating(action_kind)
3981 || data.get("revert_of").and_then(Value::as_str).is_some()
3982 {
3983 Ok(())
3984 } else {
3985 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
3986 }
3987 }
3988 }
3989}
3990
3991#[cfg(test)]
3992mod definition_body_tests {
3993 use super::safe_definition_body;
3994
3995 #[test]
3996 fn ordinary_bodies_pass() {
3997 assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
3998 assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
3999 }
4000
4001 #[test]
4002 fn a_body_cannot_close_its_own_block_or_carry_destruction() {
4003 assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
4008 assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
4009 assert!(!safe_definition_body("RECALL facts FORGET abc"));
4011 assert!(!safe_definition_body("recall facts purge older than 1d"));
4012 assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
4013 assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
4015 assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
4017 }
4018}
4019
4020#[cfg(test)]
4021mod plan_edit_tests {
4022 use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
4023 use serde_json::json;
4024
4025 fn plan() -> serde_json::Value {
4026 json!({
4027 "nodes": ["fetch", "review", "post"],
4028 "edges": [
4029 {"src": "fetch", "dst": "review"},
4030 {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
4031 ],
4032 "bindings": {"fetch": "sha256:tool1"},
4033 "retries": {"fetch": 1}
4034 })
4035 }
4036
4037 #[test]
4038 fn the_allowlist_admits_thresholds_and_refuses_topology() {
4039 assert!(plan_edit_allowed("edges.1.cond"));
4040 assert!(plan_edit_allowed("edges.1.max_cycles"));
4041 assert!(plan_edit_allowed("retries.fetch"));
4042 for path in [
4045 "nodes",
4046 "nodes.0",
4047 "edges.0.src",
4048 "edges.0.dst",
4049 "edges",
4050 "bindings.fetch",
4051 "edges.x.cond",
4052 "",
4053 ] {
4054 assert!(!plan_edit_allowed(path), "{path} must not be editable");
4055 }
4056 }
4057
4058 #[test]
4059 fn values_are_type_checked_against_the_field() {
4060 assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
4063 assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
4064 assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
4065 assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
4066 assert!(plan_value_ok("retries.fetch", &json!(3)));
4067 assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
4068 assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
4069 assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
4070 assert!(!plan_value_ok("edges.0.src", &json!("other")));
4071 }
4072
4073 #[test]
4074 fn get_reads_through_arrays_and_objects_and_absence_is_null() {
4075 let p = plan();
4076 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
4077 assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
4078 assert_eq!(plan_get(&p, "retries.review"), json!(null));
4081 assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
4082 assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
4083 }
4084
4085 #[test]
4086 fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
4087 let mut p = plan();
4088 assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
4089 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
4090 assert!(plan_set(&mut p, "retries.review", json!(2)));
4091 assert_eq!(plan_get(&p, "retries.review"), json!(2));
4092 assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
4093 assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
4094 }
4095}
4096
4097#[cfg(test)]
4098mod definition_proposal_tests {
4099 use super::is_definition_statement;
4100
4101 #[test]
4102 fn definition_statements_are_recognized_in_both_spellings() {
4103 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
4104 assert!(is_definition_statement(" define template foo AS { x }"));
4105 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
4106 assert!(!is_definition_statement("ADD fact {}"));
4108 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
4109 assert!(!is_definition_statement("FORGET abc"));
4110 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
4114 }
4115}