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 #[serde(default, skip_serializing_if = "is_zero")]
188 pub dropped_near_duplicate: u64,
189}
190
191fn is_zero(n: &u64) -> bool {
192 *n == 0
193}
194
195impl RunResult {
196 fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
197 RunResult {
198 outcome: RunOutcome::Skipped,
199 skip_reason: Some(reason),
200 new_grains,
201 new_error_events,
202 proposed: 0,
203 deduped: 0,
204 stored: 0,
205 auto_applied: 0,
206 llm_funnel: None,
207 analyzers_run: vec![],
208 analyzers_skipped: vec![],
209 }
210 }
211
212 pub fn ran(&self) -> bool {
213 self.outcome == RunOutcome::Ran
214 }
215}
216
217pub struct Engine {
220 analyzers: Vec<Box<dyn Analyzer>>,
221 policy: crate::policy::Policy,
222 llm: Option<Box<dyn crate::llm::LlmBackend>>,
225 ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
230}
231
232pub(crate) struct AnalysisPass {
233 pub(crate) survivors: Vec<Recommendation>,
234 proposed: u64,
235 deduped: u64,
236 analyzers_run: Vec<String>,
237 pub(crate) analyzers_skipped: Vec<AnalyzerSkip>,
238 llm_funnel: Option<LlmFunnel>,
239}
240
241impl Engine {
242 pub fn with_builtins() -> Self {
245 Engine {
246 analyzers: crate::analyzer::builtin_analyzers(),
247 policy: crate::policy::Policy::default(),
248 llm: None,
249 ground_llm: None,
250 }
251 }
252
253 pub fn empty() -> Self {
255 Engine {
256 analyzers: vec![],
257 policy: crate::policy::Policy::default(),
258 llm: None,
259 ground_llm: None,
260 }
261 }
262
263 pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
265 self.policy = policy;
266 self
267 }
268
269 pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
274 self.llm = Some(backend);
275 self
276 }
277
278 pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
282 self.ground_llm = Some(backend);
283 self
284 }
285
286 pub fn policy(&self) -> &crate::policy::Policy {
287 &self.policy
288 }
289
290 pub fn has_llm(&self) -> bool {
292 self.llm.is_some()
293 }
294
295 pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
297 self.analyzers.push(analyzer);
298 }
299
300 pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
301 &self.analyzers
302 }
303
304 pub fn analyze_only<S: OmsSubstrate>(
313 &self,
314 sub: &S,
315 opts: &RunOptions,
316 overrides: &BTreeMap<String, Map<String, Value>>,
317 now_ms: i64,
318 ) -> Result<Vec<Recommendation>> {
319 let persisted = LoopPersisted::from_value(sub.load_state()?)?;
320 let analysis_watermark = if opts.full_sweep {
321 None
322 } else {
323 persisted.state.watermark_ms
324 };
325 Ok(self
326 .analysis_pass(
327 sub,
328 &persisted,
329 opts,
330 overrides,
331 analysis_watermark,
332 now_ms,
333 &[],
334 )?
335 .survivors)
336 }
337
338 pub fn run<S: OmsSubstrate>(
341 &self,
342 sub: &mut S,
343 opts: &RunOptions,
344 now_ms: i64,
345 ) -> Result<RunResult> {
346 let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
347 let watermark = persisted.state.watermark_ms;
348 let analysis_watermark = if opts.full_sweep { None } else { watermark };
353
354 let new = count_new(sub, watermark)?;
355 let (new_grains, new_error_events) = (new.grains, new.error_events);
356 if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
357 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
358 }
359 let flags_set = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
363 if !flags_set && !opts.full_sweep {
364 if let Some(reason) = cadence_gate(&self.policy.cadence, &persisted, new, now_ms) {
365 return Ok(RunResult::skipped(reason, new_grains, new_error_events));
366 }
367 }
368
369 let mut outcome_inputs = measure_outcomes(sub, &mut persisted, &self.policy, now_ms)?;
374 if self.policy.premise_drift {
375 outcome_inputs.extend(detect_premise_drift(sub, &mut persisted, now_ms)?);
376 }
377
378 let AnalysisPass {
379 survivors,
380 proposed,
381 deduped,
382 analyzers_run,
383 analyzers_skipped,
384 llm_funnel,
385 } = self.analysis_pass(
386 &*sub,
387 &persisted,
388 opts,
389 &BTreeMap::new(),
390 analysis_watermark,
391 now_ms,
392 &outcome_inputs,
393 )?;
394
395 let mut stored = 0u64;
398 let mut auto_applied = 0u64;
399 for mut rec in survivors {
400 let spec = rec.to_grain_spec(LOOP_NS)?;
401 let hash = sub.put_grain(&spec)?;
402 rec.hash = hash.clone();
403 let actor = format!("engine:{}", rec.analyzer);
404 let audit = AuditRecord {
405 rec_hash: hash.clone(),
406 from: None,
407 to: RecStatus::Pending,
408 actor: actor.clone(),
409 observer_type: ObserverType::System,
410 because: "analyzer proposed".into(),
411 previous_audit_hash: None,
412 gating: None,
413 at_ms: now_ms,
414 };
415 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
416 persisted
417 .status_index
418 .insert(hash.clone(), RecStatus::Pending);
419 persisted.creators.insert(hash.clone(), actor);
420 if !matches!(rec.origin, Origin::Builtin) {
425 if let Some(trigger) = &opts.triggering_actor {
426 persisted.co_creators.insert(hash.clone(), trigger.clone());
427 }
428 }
429 persisted.audit_heads.insert(hash.clone(), audit_hash);
430 stored += 1;
431
432 if self.can_auto_apply(&*sub, &rec) {
433 self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
434 auto_applied += 1;
435 }
436 }
437
438 persisted.state.last_run_ms = Some(now_ms);
439 persisted.state.watermark_ms = Some(now_ms);
440 sub.store_state(&persisted.to_value()?)?;
441
442 Ok(RunResult {
443 outcome: RunOutcome::Ran,
444 skip_reason: None,
445 new_grains,
446 new_error_events,
447 proposed,
448 deduped,
449 stored,
450 auto_applied,
451 analyzers_run,
452 analyzers_skipped,
453 llm_funnel,
454 })
455 }
456
457 #[allow(clippy::too_many_arguments)]
461 #[allow(clippy::too_many_arguments)]
462 fn analysis_pass<S: OmsSubstrate>(
463 &self,
464 sub: &S,
465 persisted: &LoopPersisted,
466 opts: &RunOptions,
467 external_overrides: &BTreeMap<String, Map<String, Value>>,
468 analysis_watermark: Option<i64>,
469 now_ms: i64,
470 outcome_inputs: &[OutcomeInput],
471 ) -> Result<AnalysisPass> {
472 let existing = existing_dedup_keys(sub, persisted)?;
473 self.analysis_pass_inner(
474 sub,
475 persisted,
476 &self.policy,
477 opts,
478 external_overrides,
479 analysis_watermark,
480 now_ms,
481 outcome_inputs,
482 &existing,
483 None,
484 )
485 }
486
487 #[allow(clippy::too_many_arguments)]
494 pub(crate) fn analysis_pass_inner<S: OmsSubstrate>(
495 &self,
496 sub: &S,
497 persisted: &LoopPersisted,
498 policy: &crate::policy::Policy,
499 opts: &RunOptions,
500 external_overrides: &BTreeMap<String, Map<String, Value>>,
501 analysis_watermark: Option<i64>,
502 now_ms: i64,
503 outcome_inputs: &[OutcomeInput],
504 existing: &BTreeSet<String>,
505 replay: Option<&str>,
506 ) -> Result<AnalysisPass> {
507 let mut analyzers_run = Vec::new();
508 let mut analyzers_skipped = Vec::new();
509 let mut candidates: Vec<Recommendation> = Vec::new();
510 let caps = sub.capabilities();
511 let verdicts = latest_verdicts(persisted);
512
513 for analyzer in &self.analyzers {
514 let m = analyzer.manifest();
515 if let Some(why) = replay {
516 let out_of_process = m.trust_class == crate::manifest::TrustClass::Command;
517 let telemetry_fed = m.requires.contains(&crate::manifest::Capability::Telemetry);
518 if out_of_process || telemetry_fed {
519 analyzers_skipped.push(AnalyzerSkip {
520 id: m.id.clone(),
521 reason: format!(
522 "not replayed: {}",
523 if out_of_process { why } else { "telemetry rollups are not time-indexed" }
524 ),
525 });
526 continue;
527 }
528 }
529 let cfg = persisted.config.get(&m.id);
530 let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
531 if !enabled {
532 analyzers_skipped.push(AnalyzerSkip {
533 id: m.id.clone(),
534 reason: "disabled".into(),
535 });
536 continue;
537 }
538 if policy.denies(m.family()) {
539 analyzers_skipped.push(AnalyzerSkip {
540 id: m.id.clone(),
541 reason: "denied by host policy".into(),
542 });
543 continue;
544 }
545 if let Some(missing) = missing_capability(m, caps) {
546 analyzers_skipped.push(AnalyzerSkip {
547 id: m.id.clone(),
548 reason: format!("missing capability: {missing}"),
549 });
550 continue;
551 }
552 let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
553 if let Some(extra) = external_overrides.get(&m.id) {
554 for (key, value) in extra {
555 param_overrides.insert(key.clone(), value.clone());
556 }
557 }
558 let params = match m.resolve_params(¶m_overrides) {
559 Ok(p) => p,
560 Err(e) => {
561 analyzers_skipped.push(AnalyzerSkip {
562 id: m.id.clone(),
563 reason: e.to_string(),
564 });
565 continue;
566 }
567 };
568 let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
569 let ns_slice: &[String] = if ns_owned.is_empty() {
570 &opts.namespaces
571 } else {
572 &ns_owned
573 };
574 let reader: &dyn SubstrateRead = sub;
575 let ctx = AnalyzeCtx::new(
576 reader,
577 ¶ms,
578 ns_slice,
579 analysis_watermark,
580 now_ms,
581 outcome_inputs,
582 &verdicts,
583 );
584 match analyzer.analyze(&ctx) {
585 Ok(drafts) => {
586 analyzers_run.push(m.id.clone());
587 for draft in drafts {
588 match stamp(m, ¶ms, draft, now_ms) {
589 Ok(rec) => candidates.push(rec),
590 Err(e) => analyzers_skipped.push(AnalyzerSkip {
591 id: m.id.clone(),
592 reason: e.to_string(),
593 }),
594 }
595 }
596 }
597 Err(e) => analyzers_skipped.push(AnalyzerSkip {
598 id: m.id.clone(),
599 reason: e.to_string(),
600 }),
601 }
602 }
603
604 let mut funnel = LlmFunnel::default();
605 if self.llm.is_some() && replay.is_none() {
606 candidates.extend(self.discover(
607 sub,
608 &candidates,
609 analysis_watermark,
610 &opts.namespaces,
611 now_ms,
612 &mut funnel,
613 ));
614 }
615
616 let proposed = candidates.len() as u64;
617 let mut seen = BTreeSet::new();
618 let mut survivors = Vec::new();
619 for candidate in candidates {
620 let family = crate::manifest::analyzer_family(&candidate.analyzer);
621 let floor = [
622 severity_floor_for(persisted, &candidate.analyzer),
623 policy.severity_floor(family),
624 ]
625 .into_iter()
626 .flatten()
627 .max();
628 if floor.is_some_and(|floor| candidate.severity < floor) {
629 continue;
630 }
631 if !seen.insert(candidate.dedup_key.clone()) {
632 continue;
633 }
634 if existing.contains(&candidate.dedup_key) {
635 continue;
636 }
637 if persisted
638 .cooldowns
639 .get(&candidate.dedup_key)
640 .is_some_and(|until| now_ms < *until)
641 {
642 continue;
643 }
644 survivors.push(candidate);
645 }
646 let deduped = proposed - survivors.len() as u64;
647 if self.llm.is_some() && replay.is_none() {
648 self.enrich(&mut survivors);
649 }
650 Ok(AnalysisPass {
651 survivors,
652 proposed,
653 deduped,
654 analyzers_run,
655 analyzers_skipped,
656 llm_funnel: self.llm.is_some().then_some(funnel),
657 })
658 }
659
660 fn discover<S: OmsSubstrate>(
667 &self,
668 sub: &S,
669 candidates: &[Recommendation],
670 watermark: Option<i64>,
671 namespaces: &[String],
672 now_ms: i64,
673 funnel: &mut LlmFunnel,
674 ) -> Vec<Recommendation> {
675 let Some(llm) = &self.llm else {
676 return Vec::new();
677 };
678 let findings: Vec<crate::llm::FindingBrief> = candidates
679 .iter()
680 .take(32)
681 .map(|c| crate::llm::FindingBrief {
682 analyzer: c.analyzer.clone(),
683 summary: c.summary.render(),
684 target: c.target_ref.clone(),
685 severity: c.severity.as_str().to_string(),
686 })
687 .collect();
688 let attribution = self.policy.evidence_attribution;
694 let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
695 let mut bundle: BTreeSet<String> = BTreeSet::new();
696 let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
697 'cited: for c in candidates {
698 for h in &c.evidence {
699 if evidence.len() >= CITED_SEED_CAP {
708 break 'cited;
709 }
710 if !bundle.contains(h) {
711 if let Ok(Some(g)) = sub.grain(h) {
712 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
713 }
714 }
715 }
716 }
717 let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
718 vec![None]
719 } else {
720 namespaces.iter().map(|n| Some(n.as_str())).collect()
721 };
722 let opts = ReadOpts { live_only: true, since_ms: watermark };
723 let mut tool_seeded = 0usize;
744 let seed_successes = self.policy.skills.enabled || self.policy.plans.enabled;
745 'tools: for want_error in [true, false] {
746 if !want_error && !seed_successes {
747 break;
748 }
749 for ns in &scan_ns {
750 if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
751 for g in recent {
752 if tool_seeded >= TOOL_SEED_CAP
753 || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
754 {
755 break 'tools;
756 }
757 if g.is_error() != want_error {
758 continue;
759 }
760 let before = evidence.len();
761 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
762 if evidence.len() > before {
763 tool_seeded += 1;
764 }
765 }
766 }
767 }
768 }
769 if let Ok(rows) =
778 sub.grains_of_type(crate::model::grain_type::OBSERVATION, Some(crate::eval::HARNESS_NS), opts)
779 {
780 let mut seeded = 0usize;
781 for g in rows {
782 if seeded >= HARNESS_SEED_CAP || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE {
783 break;
784 }
785 let kind = g.str_field("observation_kind").unwrap_or_default();
786 if !HARNESS_EVIDENCE_KINDS.contains(&kind) {
787 continue;
788 }
789 let before = evidence.len();
790 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
791 if evidence.len() > before {
792 seeded += 1;
793 }
794 }
795 }
796 'notes: for ns in &scan_ns {
806 if let Ok(recent) =
807 sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
808 {
809 for g in recent {
810 if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
811 break 'notes;
812 }
813 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
814 }
815 }
816 }
817 'seed: for gt in [
818 crate::model::grain_type::FACT,
819 crate::model::grain_type::OBSERVATION,
820 ] {
821 for ns in &scan_ns {
822 if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
823 for g in recent {
824 if evidence.len() >= EVIDENCE_CAP {
825 break 'seed;
826 }
827 push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
828 }
829 }
830 }
831 }
832 funnel.evidence = evidence.len() as u64;
833 if evidence.is_empty() {
834 return Vec::new(); }
836 let (approved, rejected) = self.llm_history(sub);
841 let base = match self.policy.discover_objective {
842 crate::policy::DiscoverObjective::ReviewQueue => DISCOVER_INSTRUCTIONS,
843 crate::policy::DiscoverObjective::Learner => DISCOVER_LEARNER_INSTRUCTIONS,
844 };
845 let mut instructions = base.to_string();
848 if self.policy.skills.enabled {
849 instructions.push_str(&skill_instructions(self.policy.skills.min_steps));
850 }
851 if self.policy.plans.enabled {
852 instructions.push_str(&plan_instructions(self.policy.plans.min_nodes));
853 }
854 if findings.iter().any(|f| f.analyzer.starts_with("loop.lesson_pile/")) {
857 instructions.push_str(CONSOLIDATION_INSTRUCTIONS);
858 }
859 let request = crate::llm::LlmRequest {
860 loop_proto: 1,
861 op: "discover",
862 instructions: &instructions,
863 findings: findings.clone(),
864 evidence: evidence.clone(),
865 rejected,
866 approved,
867 };
868 let Ok(body) = serde_json::to_string(&request) else {
869 return Vec::new();
870 };
871 let raw = match llm.complete(&body) {
872 Ok(r) => r,
873 Err(_) => return Vec::new(), };
875 let caps = sub.capabilities();
879 let mut validated: Vec<ValidatedDraft> = Vec::new();
880 let drafts: Vec<_> = crate::llm::parse_discover(&raw)
881 .recommendations
882 .into_iter()
883 .take(crate::llm::MAX_LLM_DRAFTS)
884 .collect();
885 funnel.proposed = drafts.len() as u64;
886 let id_to_hash: std::collections::BTreeMap<&str, &str> = evidence
887 .iter()
888 .map(|e| (e.id.as_str(), e.hash.as_str()))
889 .collect();
890 for d in drafts {
891 let mut cited: Vec<String> = Vec::new();
892 for c in &d.evidence {
893 if let Some(h) = resolve_citation(c, &bundle, &id_to_hash) {
894 if !cited.contains(&h) {
895 cited.push(h);
896 }
897 }
898 }
899 if cited.is_empty() {
900 funnel.dropped_uncited += 1;
901 continue; }
903 let Ok(target) = TargetRef::parse(&d.target) else {
904 funnel.dropped_target += 1;
905 continue;
906 };
907 let tc = target.target_class();
908 if !matches!(tc, "memory" | "query" | "code") {
914 funnel.dropped_target += 1;
915 continue;
916 }
917 let thin = cited.len() < self.policy.min_evidence as usize;
924 if thin {
925 funnel.advisory_thin_evidence += 1;
926 }
927 let resolved = if thin {
928 None
929 } else {
930 resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps, &self.policy)
931 };
932 if tc == "code" && resolved.is_none() {
937 funnel.dropped_target += 1;
938 continue;
939 }
940 validated.push(ValidatedDraft {
941 draft: d,
942 target_ref: target.as_string(),
943 cited,
944 resolved,
945 });
946 }
947 funnel.cited = validated.len() as u64;
948 if validated.is_empty() {
949 return Vec::new();
950 }
951 let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
958 let outcome_metric = self.outcome_metric_template(sub);
959 self.verify_drafts(sub, &**llm, ground, validated, &evidence, outcome_metric, now_ms, funnel)
960 }
961
962 fn outcome_metric_template<S: OmsSubstrate>(
970 &self,
971 sub: &S,
972 ) -> Option<crate::recommendation::MetricSnapshot> {
973 let e = self.policy.outcome_evalset.as_ref()?;
974 let run = crate::eval::newest_eval_run(sub, &e.hash, None).ok().flatten()?;
975 let baseline = crate::eval::run_value(&run, &e.field)?;
976 let schedule = e.schedule();
980 let ms_only: Vec<i64> = schedule.iter().filter_map(Checkpoint::as_ms).collect();
981 let all_ms = ms_only.len() == schedule.len();
982 let horizons = if ms_only.is_empty() { vec![86_400_000] } else { ms_only };
983 Some(crate::recommendation::MetricSnapshot {
984 metric: format!("evalset:{}:{}", e.hash, e.field),
985 baseline,
986 unit: e.field.clone(),
987 n: run.total(),
988 window: "per-run".into(),
989 subject: None,
990 namespace: None,
991 relation: None,
992 query: format!(
993 "RECALL facts WHERE subject = \"evalset:{}\" AND relation = \"mg:eval_run\"",
994 e.hash
995 ),
996 review_after_ms: horizons[0],
997 horizons_ms: if all_ms { horizons } else { Vec::new() },
998 checkpoints: if all_ms { Vec::new() } else { schedule },
999 higher_is_better: e.higher_is_better,
1000 })
1001 }
1002
1003 #[allow(clippy::too_many_arguments)]
1010 #[allow(clippy::too_many_arguments)]
1011 fn verify_drafts<S: SubstrateRead>(
1012 &self,
1013 sub: &S,
1014 llm: &dyn crate::llm::LlmBackend,
1015 ground: &dyn crate::llm::LlmBackend,
1016 validated: Vec<ValidatedDraft>,
1017 evidence: &[crate::llm::EvidenceItem],
1018 outcome_metric: Option<crate::recommendation::MetricSnapshot>,
1019 now_ms: i64,
1020 funnel: &mut LlmFunnel,
1021 ) -> Vec<Recommendation> {
1022 use crate::llm::*;
1023 let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
1024 evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
1025 let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
1026 cited
1027 .iter()
1028 .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
1029 .collect()
1030 };
1031
1032 let claims: Vec<GroundItem> = validated
1036 .iter()
1037 .enumerate()
1038 .map(|(i, v)| GroundItem {
1039 id: i,
1040 claim: claim_text(&v.draft, v.resolved.as_ref()),
1041 evidence: ev_for(&v.cited),
1042 })
1043 .collect();
1044 let ground_req = GroundRequest {
1045 loop_proto: 1,
1046 op: "ground",
1047 instructions: GROUND_INSTRUCTIONS,
1048 claims,
1049 };
1050 let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
1056 .ok()
1057 .and_then(|b| ground.complete(&b).ok())
1058 {
1059 Some(raw) => {
1060 let parsed = parse_ground(&raw);
1061 funnel.ground_verdicts = parsed.results.len() as u64;
1062 parsed
1063 .results
1064 .into_iter()
1065 .filter(|r| r.supported)
1066 .map(|r| r.id)
1067 .collect()
1068 }
1069 None => {
1070 funnel.ground_call_failed = true;
1071 return Vec::new();
1072 }
1073 };
1074 funnel.grounded = grounded.len() as u64;
1075 if grounded.is_empty() {
1076 return Vec::new();
1077 }
1078
1079 let items: Vec<VerifyItem> = validated
1085 .iter()
1086 .enumerate()
1087 .filter(|(i, _)| grounded.contains(i))
1088 .map(|(i, v)| VerifyItem {
1089 id: i,
1090 summary: claim_text(&v.draft, v.resolved.as_ref()),
1092 target: v.target_ref.clone(),
1093 evidence: ev_for(&v.cited),
1094 })
1095 .collect();
1096 let verify_req = VerifyRequest {
1097 loop_proto: 1,
1098 op: "verify",
1099 instructions: VERIFY_INSTRUCTIONS,
1100 findings: items,
1101 };
1102 let verdicts: std::collections::BTreeMap<usize, f64> =
1103 match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
1104 Some(raw) => parse_verify(&raw)
1105 .results
1106 .into_iter()
1107 .filter(|r| r.keep)
1108 .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
1109 .collect(),
1110 None => return Vec::new(),
1111 };
1112
1113 funnel.kept = verdicts.len() as u64;
1114 let mut out = Vec::new();
1118 for (i, v) in validated.into_iter().enumerate() {
1119 if let Some(&conf) = verdicts.get(&i) {
1120 if conf >= MIN_LLM_CONFIDENCE {
1121 let near = match v.resolved.as_ref() {
1130 Some(r) if r.action != ActionKind::Consolidate => r
1131 .fact_fields
1132 .as_ref()
1133 .filter(|f| f.get("relation").and_then(Value::as_str) == Some("lesson"))
1134 .map(|f| {
1135 near_duplicates_of(
1136 sub,
1137 f.get("subject").and_then(Value::as_str).unwrap_or(""),
1138 f.get("namespace").and_then(Value::as_str),
1139 f.get("object").and_then(Value::as_str).unwrap_or(""),
1140 )
1141 })
1142 .unwrap_or_default(),
1143 _ => Vec::new(),
1144 };
1145 if !near.is_empty()
1146 && self.policy.near_duplicate == crate::policy::NearDuplicateMode::Suppress
1147 {
1148 funnel.dropped_near_duplicate += 1;
1149 continue;
1150 }
1151 let mut rec = stamp_llm(
1152 llm.model(),
1153 &v.draft,
1154 v.target_ref,
1155 v.cited,
1156 v.resolved,
1157 conf,
1158 now_ms,
1159 );
1160 if let Some(best) = near.first() {
1161 rec.summary.args.insert("near_count".into(), Value::from(near.len() as u64));
1163 rec.summary.args.insert("near_score".into(), Value::from(best.score));
1164 rec.summary.args.insert("near_method".into(), Value::from(best.method.clone()));
1165 rec.summary.args.insert(
1166 "near_hash".into(),
1167 Value::from(best.hash.chars().take(12).collect::<String>()),
1168 );
1169 rec.summary.template_id = "llm.lesson_near_duplicate".into();
1170 rec.near_duplicate_of = near;
1171 }
1172 if rec.rollbackable {
1176 rec.metric = outcome_metric.clone();
1177 }
1178 out.push(rec);
1179 }
1180 }
1181 }
1182 funnel.stored = out.len() as u64;
1183 out
1184 }
1185
1186 fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
1191 const MAX: usize = 20;
1192 let Ok(mut recs) = self.recommendations(sub, None) else {
1193 return (Vec::new(), Vec::new());
1194 };
1195 recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
1196 recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
1197 let mut approved = Vec::new();
1198 let mut rejected = Vec::new();
1199 for r in &recs {
1200 match r.status {
1201 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
1202 if approved.len() < MAX =>
1203 {
1204 approved.push(r.summary.render());
1205 }
1206 RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
1207 _ => {}
1208 }
1209 }
1210 (approved, rejected)
1211 }
1212
1213 fn enrich(&self, survivors: &mut [Recommendation]) {
1218 let Some(llm) = &self.llm else {
1219 return;
1220 };
1221 if survivors.is_empty() {
1222 return;
1223 }
1224 let findings: Vec<crate::llm::FindingBrief> = survivors
1225 .iter()
1226 .map(|r| crate::llm::FindingBrief {
1227 analyzer: r.analyzer.clone(),
1228 summary: r.summary.render(),
1229 target: r.target_ref.clone(),
1230 severity: r.severity.as_str().to_string(),
1231 })
1232 .collect();
1233 let request = crate::llm::LlmRequest {
1234 loop_proto: 1,
1235 op: "enrich",
1236 instructions: ENRICH_INSTRUCTIONS,
1237 findings,
1238 evidence: Vec::new(),
1239 rejected: Vec::new(),
1240 approved: Vec::new(),
1241 };
1242 let Ok(body) = serde_json::to_string(&request) else {
1243 return;
1244 };
1245 let raw = match llm.complete(&body) {
1246 Ok(r) => r,
1247 Err(_) => return,
1248 };
1249 for note in crate::llm::parse_enrich(&raw).notes {
1250 if note.guidance.trim().is_empty() {
1251 continue;
1252 }
1253 if let Some(r) = survivors
1254 .iter_mut()
1255 .find(|r| r.target_ref == note.target && r.guidance.is_none())
1256 {
1257 r.guidance = Some(crate::llm::cap(¬e.guidance, crate::llm::MAX_GUIDANCE_LEN));
1258 }
1259 }
1260 }
1261
1262 fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
1271 if !rec.origin.auto_apply_eligible() || rec.destructive {
1272 return false;
1273 }
1274 let manifest_ok = self
1278 .analyzers
1279 .iter()
1280 .map(|a| a.manifest())
1281 .find(|m| m.id == rec.analyzer)
1282 .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
1283 if !manifest_ok {
1284 return false;
1285 }
1286 let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1287 return false;
1288 };
1289 let family = crate::manifest::analyzer_family(&rec.analyzer);
1290 if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1291 return false;
1292 }
1293 match &rec.proposal {
1298 Proposal::Cal { cal } => cal
1299 .lines()
1300 .map(str::trim)
1301 .filter(|l| !l.is_empty())
1302 .all(|l| supersede_is_value_identical(sub, l)),
1303 _ => false,
1304 }
1305 }
1306
1307 fn auto_apply<S: OmsSubstrate>(
1310 &self,
1311 sub: &mut S,
1312 p: &mut LoopPersisted,
1313 rec: &Recommendation,
1314 now_ms: i64,
1315 ) -> Result<()> {
1316 let mut created = Vec::new();
1317 if let Proposal::Cal { cal } = &rec.proposal {
1318 if cal.lines().map(str::trim).any(is_definition_statement) {
1323 return Err(Error::InvalidProposal(
1324 "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1325 auto-applied: it changes what every future context contains, so it \
1326 requires a human APPROVE + APPLY with BECAUSE"
1327 .into(),
1328 ));
1329 }
1330 for r in sub.execute_cal(cal)? {
1331 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1332 created.push(h.to_string());
1333 }
1334 }
1335 }
1336 let applied = AppliedRecord {
1337 applied_at_ms: now_ms,
1338 target_ref: rec.target_ref.clone(),
1339 rollbackable: rec.rollbackable,
1340 created_hashes: created,
1341 inverse_cal: None,
1342 metric: rec.metric.clone(),
1343 };
1344 let prev = p.audit_heads.get(&rec.hash).cloned();
1345 let audit = AuditRecord {
1346 rec_hash: rec.hash.clone(),
1347 from: Some(RecStatus::Pending),
1348 to: RecStatus::Applied,
1349 actor: "policy:auto".into(),
1350 observer_type: ObserverType::Policy,
1351 because: "auto-applied per host policy".into(),
1352 previous_audit_hash: prev,
1353 gating: None,
1354 at_ms: now_ms,
1355 };
1356 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1357 p.audit_heads.insert(rec.hash.clone(), audit_hash);
1358 p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1359 p.applied.insert(rec.hash.clone(), applied);
1360 Ok(())
1361 }
1362
1363 #[allow(clippy::too_many_arguments)]
1366 pub fn review<S: OmsSubstrate>(
1367 &self,
1368 sub: &mut S,
1369 rec_hash: &str,
1370 decision: Decision,
1371 actor: &str,
1372 observer: ObserverType,
1373 scopes: &ScopeSet,
1374 because: &str,
1375 now_ms: i64,
1376 ) -> Result<()> {
1377 if !scopes.has(Scope::Review) {
1378 return Err(Error::ScopeDenied("review".into()));
1379 }
1380 let because = validate_because(because)?;
1381 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1382 let status = *p
1383 .status_index
1384 .get(rec_hash)
1385 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1386 let to = match decision {
1387 Decision::Approve => RecStatus::Approved,
1388 Decision::Reject => RecStatus::Rejected,
1389 };
1390 if !status.can_transition_to(to, false) {
1391 return Err(Error::LifecycleViolation(format!(
1392 "{} -> {}",
1393 status.as_str(),
1394 to.as_str()
1395 )));
1396 }
1397 if to == RecStatus::Approved {
1398 if let Some(creator) = p.creators.get(rec_hash) {
1399 if creator == actor {
1400 return Err(Error::SelfApproval(format!(
1401 "{actor} created this recommendation"
1402 )));
1403 }
1404 }
1405 if let Some(trigger) = p.co_creators.get(rec_hash) {
1406 if trigger == actor {
1407 return Err(Error::SelfApproval(format!(
1408 "{actor} triggered the run that authored this recommendation"
1409 )));
1410 }
1411 }
1412 }
1413 let prev = p.audit_heads.get(rec_hash).cloned();
1414 let audit = AuditRecord {
1415 rec_hash: rec_hash.into(),
1416 from: Some(status),
1417 to,
1418 actor: actor.into(),
1419 observer_type: observer,
1420 because,
1421 previous_audit_hash: prev,
1422 gating: None,
1423 at_ms: now_ms,
1424 };
1425 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1426 p.audit_heads.insert(rec_hash.into(), audit_hash);
1427 p.status_index.insert(rec_hash.into(), to);
1428 if to == RecStatus::Rejected {
1429 if let Ok(rec) = load_rec(sub, rec_hash) {
1430 strike_cooldown(&mut p, rec.dedup_key, now_ms);
1431 }
1432 }
1433 sub.store_state(&p.to_value()?)?;
1434 Ok(())
1435 }
1436
1437 pub fn preflight_apply<S: OmsSubstrate>(
1454 &self,
1455 sub: &S,
1456 rec_hash: &str,
1457 scopes: &ScopeSet,
1458 allow_destructive: bool,
1459 has_gating: bool,
1460 ) -> Result<()> {
1461 if !scopes.has(Scope::Apply) {
1462 return Err(Error::ScopeDenied("apply".into()));
1463 }
1464 let rec = load_rec(sub, rec_hash)?;
1465 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1466 return Err(Error::DestructiveGated(
1467 "destructive apply requires admin scope + allow_destructive".into(),
1468 ));
1469 }
1470 ensure_executable(rec.action_kind, &rec.proposal)?;
1471 if requires_gating(rec.action_kind) && !has_gating {
1472 return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1473 }
1474 Ok(())
1475 }
1476
1477 #[allow(clippy::too_many_arguments)]
1481 pub fn apply<S: OmsSubstrate>(
1482 &self,
1483 sub: &mut S,
1484 rec_hash: &str,
1485 actor: &str,
1486 observer: ObserverType,
1487 scopes: &ScopeSet,
1488 because: &str,
1489 allow_destructive: bool,
1490 now_ms: i64,
1491 ) -> Result<AppliedRecord> {
1492 self.apply_inner(
1493 sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1494 )
1495 }
1496
1497 pub fn gating_evidence<S: OmsSubstrate>(
1504 &self,
1505 sub: &S,
1506 rec_hash: &str,
1507 run_id: &str,
1508 ) -> Result<crate::recommendation::GatingEvidence> {
1509 let rec = self
1510 .recommendations(sub, None)?
1511 .into_iter()
1512 .find(|r| r.hash == rec_hash)
1513 .ok_or_else(|| {
1514 Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1515 })?;
1516 let pin = rec.evalset_hash.ok_or_else(|| {
1517 Error::InvalidProposal(
1518 "this recommendation pins no evalset — a gating run applies only \
1519 to code and adapter revisions"
1520 .into(),
1521 )
1522 })?;
1523 match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1528 Some(run) => Ok(crate::recommendation::GatingEvidence {
1529 evalset_hash: pin,
1530 run_id: run.run_id,
1531 passed: run.passed,
1532 failed: run.failed,
1533 }),
1534 None => Err(Error::InvalidProposal(format!(
1535 "no recorded gate run '{run_id}' for evalset {pin} — run \
1536 `areev eval run --evalset {pin} ...` first"
1537 ))),
1538 }
1539 }
1540
1541 #[allow(clippy::too_many_arguments)]
1545 pub fn apply_gated<S: OmsSubstrate>(
1546 &self,
1547 sub: &mut S,
1548 rec_hash: &str,
1549 actor: &str,
1550 observer: ObserverType,
1551 scopes: &ScopeSet,
1552 because: &str,
1553 allow_destructive: bool,
1554 gating: &crate::recommendation::GatingEvidence,
1555 now_ms: i64,
1556 ) -> Result<AppliedRecord> {
1557 self.apply_inner(
1558 sub,
1559 rec_hash,
1560 actor,
1561 observer,
1562 scopes,
1563 because,
1564 allow_destructive,
1565 Some(gating),
1566 now_ms,
1567 )
1568 }
1569
1570 #[allow(clippy::too_many_arguments)]
1571 fn apply_inner<S: OmsSubstrate>(
1572 &self,
1573 sub: &mut S,
1574 rec_hash: &str,
1575 actor: &str,
1576 observer: ObserverType,
1577 scopes: &ScopeSet,
1578 because: &str,
1579 allow_destructive: bool,
1580 gating: Option<&crate::recommendation::GatingEvidence>,
1581 now_ms: i64,
1582 ) -> Result<AppliedRecord> {
1583 if !scopes.has(Scope::Apply) {
1584 return Err(Error::ScopeDenied("apply".into()));
1585 }
1586 let because = validate_because(because)?;
1587 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1588 let status = *p
1589 .status_index
1590 .get(rec_hash)
1591 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1592 if !status.can_transition_to(RecStatus::Applied, false) {
1593 return Err(Error::LifecycleViolation(format!(
1594 "{} -> applied (approve first)",
1595 status.as_str()
1596 )));
1597 }
1598 let rec = load_rec(sub, rec_hash)?;
1599 if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1600 return Err(Error::DestructiveGated(
1601 "destructive apply requires admin scope + allow_destructive".into(),
1602 ));
1603 }
1604 if requires_gating(rec.action_kind) {
1609 let g = gating
1610 .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1611 let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1612 if g.evalset_hash != pin {
1613 return Err(Error::InvalidProposal(format!(
1614 "gating ran evalset {} but the recommendation is pinned \
1615 to {pin} (Rule E1)",
1616 g.evalset_hash
1617 )));
1618 }
1619 match sub.grain(pin)? {
1620 Some(evalset) if evalset.is_live() => {}
1621 Some(_) => {
1622 return Err(Error::InvalidProposal(
1623 "the pinned evalset was superseded after gating — \
1624 the recommendation must re-gate (Rule E1)"
1625 .into(),
1626 ))
1627 }
1628 None => {
1629 return Err(Error::InvalidProposal(format!(
1630 "pinned evalset {pin} not found in the substrate"
1631 )))
1632 }
1633 }
1634 if g.failed > 0 {
1635 return Err(Error::InvalidProposal(format!(
1636 "the gating run failed {}/{} cases — a failing gate \
1637 admits nothing",
1638 g.failed,
1639 g.passed + g.failed
1640 )));
1641 }
1642 }
1643
1644 let mut created = Vec::new();
1646 let mut inverse_cal: Option<String> = None;
1650 match &rec.proposal {
1651 Proposal::Cal { cal } => {
1652 for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1653 if !is_definition_statement(line) {
1654 continue;
1655 }
1656 match sub.definition_inverse(line)? {
1657 Some(inv) => inverse_cal = Some(inv),
1658 None => {
1659 return Err(Error::InvalidProposal(format!(
1660 "this substrate cannot record a rollback inverse for {line:?}; \
1661 a definition rewrite that ROLLBACK could not undo is refused \
1662 rather than applied"
1663 )))
1664 }
1665 }
1666 }
1667 let rows = sub.execute_cal(cal)?;
1668 for r in rows {
1669 if let Some(h) = r.get("hash").and_then(Value::as_str) {
1670 created.push(h.to_string());
1671 }
1672 }
1673 }
1674 Proposal::Edit { .. } => {
1677 return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1678 }
1679 Proposal::Data { data } if requires_gating(rec.action_kind) => {
1687 let relation = if rec.action_kind == ActionKind::AdapterRevision {
1688 "mg:adapter_promotion"
1689 } else {
1690 "mg:code_promotion"
1691 };
1692 let mut promoted = data.clone();
1700 if let Some(Value::String(src)) = promoted.remove("source") {
1701 let address = sub.put_blob(src.as_bytes())?;
1702 promoted.insert("code_address".into(), Value::from(address));
1703 }
1704 let mut spec = crate::substrate::GrainSpec::new(
1705 crate::model::grain_type::FACT,
1706 LOOP_NS,
1707 )
1708 .with_field("subject", rec.target_ref.clone())
1709 .with_field("relation", relation)
1710 .with_field(
1711 "object",
1712 serde_json::to_string(&promoted)
1713 .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1714 )
1715 .with_field("rec_hash", rec_hash.to_string());
1716 if let Some(g) = gating {
1717 spec = spec
1718 .with_field("gating_evalset", g.evalset_hash.clone())
1719 .with_field("gating_run_id", g.run_id.clone());
1720 }
1721 created.push(sub.put_grain(&spec)?);
1722 }
1723 Proposal::Data { data } => {
1724 let revert_of = data
1729 .get("revert_of")
1730 .and_then(Value::as_str)
1731 .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1732 self.rollback(
1733 sub,
1734 revert_of,
1735 actor,
1736 observer,
1737 scopes,
1738 &because,
1739 now_ms,
1740 )?;
1741 p = LoopPersisted::from_value(sub.load_state()?)?;
1744 if let Ok(reverted) = load_rec(sub, revert_of) {
1754 strike_cooldown(&mut p, reverted.dedup_key, now_ms);
1755 }
1756 }
1757 }
1758
1759 let applied = AppliedRecord {
1760 applied_at_ms: now_ms,
1761 target_ref: rec.target_ref.clone(),
1762 rollbackable: rec.rollbackable,
1763 created_hashes: created,
1764 inverse_cal,
1765 metric: rec.metric.clone(),
1766 };
1767 let prev = p.audit_heads.get(rec_hash).cloned();
1768 let audit = AuditRecord {
1769 rec_hash: rec_hash.into(),
1770 from: Some(status),
1771 to: RecStatus::Applied,
1772 actor: actor.into(),
1773 observer_type: observer,
1774 because,
1775 previous_audit_hash: prev,
1776 gating: gating.cloned(),
1777 at_ms: now_ms,
1778 };
1779 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1780 p.audit_heads.insert(rec_hash.into(), audit_hash);
1781 p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1782 p.applied.insert(rec_hash.into(), applied.clone());
1783 sub.store_state(&p.to_value()?)?;
1784 Ok(applied)
1785 }
1786
1787 #[allow(clippy::too_many_arguments)]
1790 pub fn rollback<S: OmsSubstrate>(
1791 &self,
1792 sub: &mut S,
1793 rec_hash: &str,
1794 actor: &str,
1795 observer: ObserverType,
1796 scopes: &ScopeSet,
1797 because: &str,
1798 now_ms: i64,
1799 ) -> Result<()> {
1800 if !scopes.has(Scope::Apply) {
1801 return Err(Error::ScopeDenied("apply".into()));
1802 }
1803 let because = validate_because(because)?;
1804 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1805 let status = *p
1806 .status_index
1807 .get(rec_hash)
1808 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1809 if !status.can_transition_to(RecStatus::RolledBack, false) {
1810 return Err(Error::LifecycleViolation(format!(
1811 "{} -> rolled_back",
1812 status.as_str()
1813 )));
1814 }
1815 let applied = p
1816 .applied
1817 .get(rec_hash)
1818 .cloned()
1819 .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1820 if !applied.rollbackable {
1821 return Err(Error::LifecycleViolation(
1822 "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1823 ));
1824 }
1825 for h in &applied.created_hashes {
1826 sub.retract(h, &format!("rollback of {rec_hash}"))?;
1827 }
1828 if let Some(inverse) = &applied.inverse_cal {
1835 sub.execute_cal(inverse)?;
1836 }
1837 let prev = p.audit_heads.get(rec_hash).cloned();
1838 let audit = AuditRecord {
1839 rec_hash: rec_hash.into(),
1840 from: Some(status),
1841 to: RecStatus::RolledBack,
1842 actor: actor.into(),
1843 observer_type: observer,
1844 because,
1845 previous_audit_hash: prev,
1846 gating: None,
1847 at_ms: now_ms,
1848 };
1849 let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1850 p.audit_heads.insert(rec_hash.into(), audit_hash);
1851 p.status_index
1852 .insert(rec_hash.into(), RecStatus::RolledBack);
1853 sub.store_state(&p.to_value()?)?;
1854 Ok(())
1855 }
1856
1857 pub fn recommendations<S: OmsSubstrate>(
1862 &self,
1863 sub: &S,
1864 status_filter: Option<RecStatus>,
1865 ) -> Result<Vec<Recommendation>> {
1866 let p = LoopPersisted::from_value(sub.load_state()?)?;
1867 let grains = sub.grains_of_type(
1868 crate::model::grain_type::RECOMMENDATION,
1869 Some(LOOP_NS),
1870 ReadOpts {
1871 live_only: false,
1872 since_ms: None,
1873 },
1874 )?;
1875 let mut out = Vec::new();
1876 for g in grains {
1877 let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1878 rec.status = p
1879 .status_index
1880 .get(&g.hash)
1881 .copied()
1882 .unwrap_or(RecStatus::Pending);
1883 if let Some(f) = status_filter {
1884 if rec.status != f {
1885 continue;
1886 }
1887 }
1888 out.push(rec);
1889 }
1890 out.sort_by(|a, b| {
1899 b.severity
1900 .cmp(&a.severity)
1901 .then(a.created_at_ms.cmp(&b.created_at_ms))
1902 .then(a.dedup_key.cmp(&b.dedup_key))
1903 .then(a.hash.cmp(&b.hash))
1904 });
1905 Ok(out)
1906 }
1907
1908 pub fn analyzer_settings<S: OmsSubstrate>(
1911 &self,
1912 sub: &S,
1913 ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1914 let p = LoopPersisted::from_value(sub.load_state()?)?;
1915 Ok(self
1916 .analyzers
1917 .iter()
1918 .map(|a| {
1919 let m = a.manifest();
1920 let cfg = p.config.get(&m.id);
1921 crate::config::AnalyzerSetting {
1922 id: m.id.clone(),
1923 title: m.title.clone(),
1924 description: m.description.clone(),
1925 tier: format!("{:?}", m.tier),
1926 trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1927 default_on: m.default_on,
1928 enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1929 severity_floor: cfg
1930 .and_then(|c| c.severity_floor)
1931 .map(|s| s.as_str().to_string()),
1932 }
1933 })
1934 .collect())
1935 }
1936
1937 pub fn set_analyzer_config<S: OmsSubstrate>(
1943 &self,
1944 sub: &mut S,
1945 analyzer_id: &str,
1946 update: crate::config::AnalyzerConfigUpdate,
1947 scopes: &ScopeSet,
1948 ) -> Result<crate::config::AnalyzerConfig> {
1949 if !scopes.has(Scope::Admin) {
1950 return Err(Error::ScopeDenied("admin".into()));
1951 }
1952 let manifest = self
1953 .analyzers
1954 .iter()
1955 .map(|a| a.manifest())
1956 .find(|m| m.id == analyzer_id)
1957 .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1958 if let Some(params) = &update.params {
1960 manifest.resolve_params(params)?;
1961 }
1962 let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1963 let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1964 if let Some(enabled) = update.enabled {
1965 cfg.enabled = Some(enabled);
1966 }
1967 if update.clear_floor {
1968 cfg.severity_floor = None;
1969 } else if let Some(floor) = update.severity_floor {
1970 cfg.severity_floor = Some(floor);
1971 }
1972 if let Some(params) = update.params {
1973 cfg.params = params;
1974 }
1975 if let Some(ns) = update.namespaces {
1976 cfg.namespaces = ns;
1977 }
1978 let stored = cfg.clone();
1979 sub.store_state(&p.to_value()?)?;
1980 Ok(stored)
1981 }
1982
1983 pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1986 let p = LoopPersisted::from_value(sub.load_state()?)?;
1987 let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1988 out.sort_by(|a, b| {
1992 a.measured_at_ms
1993 .cmp(&b.measured_at_ms)
1994 .then(a.horizon_ms.cmp(&b.horizon_ms))
1995 .then(a.metric.cmp(&b.metric))
1996 .then(a.rec_hash.cmp(&b.rec_hash))
1997 });
1998 Ok(out)
1999 }
2000
2001 pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
2005 let p = LoopPersisted::from_value(sub.load_state()?)?;
2006 let new = count_new(sub, p.state.watermark_ms)?;
2007 let (grains_since_run, error_events_since_run) = (new.grains, new.error_events);
2008 let recs = self.recommendations(sub, None)?;
2009 let mut pending = 0;
2010 let mut applied = 0;
2011 for r in &recs {
2012 match r.status {
2013 RecStatus::Pending => pending += 1,
2014 RecStatus::Applied => applied += 1,
2015 _ => {}
2016 }
2017 }
2018 let stale = match p.state.last_run_ms {
2020 None => true,
2021 Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
2022 };
2023 Ok(Health {
2024 last_run_ms: p.state.last_run_ms,
2025 grains_since_run,
2026 error_events_since_run,
2027 pending,
2028 applied,
2029 total: recs.len() as u64,
2030 stale,
2031 })
2032 }
2033
2034 pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
2039 let recs = self.recommendations(sub, None)?;
2040 let mut m = LlmMetrics::default();
2041 for r in &recs {
2042 if !matches!(r.origin, Origin::Llm { .. }) {
2043 continue;
2044 }
2045 m.proposed += 1;
2046 match r.status {
2047 RecStatus::Pending => m.pending += 1,
2048 RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
2049 RecStatus::Rejected => m.rejected += 1,
2050 RecStatus::Expired => {}
2051 }
2052 }
2053 let decided = m.approved + m.rejected;
2054 m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
2055 Ok(m)
2056 }
2057}
2058
2059#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
2061pub struct Health {
2062 #[serde(skip_serializing_if = "Option::is_none")]
2063 pub last_run_ms: Option<i64>,
2064 pub grains_since_run: u64,
2065 pub error_events_since_run: u64,
2066 pub pending: u64,
2067 pub applied: u64,
2068 pub total: u64,
2069 pub stale: bool,
2072}
2073
2074#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
2076pub struct LlmMetrics {
2077 pub proposed: u64,
2080 pub pending: u64,
2081 pub approved: u64,
2083 pub rejected: u64,
2084 #[serde(skip_serializing_if = "Option::is_none")]
2086 pub approval_rate: Option<f64>,
2087}
2088
2089fn measure_outcomes<S: OmsSubstrate>(
2096 sub: &S,
2097 p: &mut LoopPersisted,
2098 policy: &crate::policy::Policy,
2099 now_ms: i64,
2100) -> Result<Vec<OutcomeInput>> {
2101 let mut due: Vec<(String, crate::config::AppliedRecord, Checkpoint)> = Vec::new();
2105 for (h, a) in &p.applied {
2106 if p.status_index.get(h) != Some(&RecStatus::Applied) {
2107 continue;
2108 }
2109 let Some(metric) = &a.metric else { continue };
2110 let done = p.measured.get(h).cloned().unwrap_or_default();
2111 for cp in metric.schedule() {
2112 if !done.contains(&cp) && checkpoint_due(sub, metric, a.applied_at_ms, cp, now_ms)? {
2113 due.push((h.clone(), a.clone(), cp));
2114 }
2115 }
2116 }
2117
2118 let mut out = Vec::new();
2119 for (rec_hash, applied, checkpoint) in due {
2120 let metric = applied.metric.as_ref().unwrap();
2121 let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
2122 continue; };
2124 let bound = cost_bound_for(policy, metric);
2125 let base = baseline_at_apply(
2126 sub,
2127 metric,
2128 applied.applied_at_ms,
2129 baseline_kind_for(policy, metric),
2130 bound.map(|b| b.field.as_str()),
2131 )?;
2132 let baseline = base.value;
2133 let tolerance = tolerance_for(policy, metric, &base);
2136 let regressed = crate::recommendation::is_regression(
2137 baseline,
2138 current,
2139 metric.higher_is_better,
2140 tolerance,
2141 );
2142 let (current_run_id, cost) = current_run_and_cost(sub, metric, applied.applied_at_ms, bound, &base)?;
2147 let costlier = cost.as_ref().is_some_and(|c| c.breached());
2148 let verdict = if regressed {
2149 "regressed"
2150 } else if costlier {
2151 "held_costlier"
2152 } else {
2153 "held"
2154 };
2155 p.outcomes.entry(rec_hash.clone()).or_default().push(
2156 crate::recommendation::OutcomeResult {
2157 rec_hash: rec_hash.clone(),
2158 metric: metric.metric.clone(),
2159 baseline,
2160 current,
2161 verdict: verdict.into(),
2162 baseline_kind: base.kind.into(),
2163 baseline_run_id: base.run_id.clone(),
2164 best_before: base.best_before,
2165 tolerance,
2166 current_run_id: current_run_id.clone(),
2167 cost: cost.clone(),
2168 horizon_ms: checkpoint.as_ms().unwrap_or(0),
2171 checkpoint: checkpoint.as_ms().is_none().then_some(checkpoint),
2172 measured_at_ms: now_ms,
2173 },
2174 );
2175 p.measured.entry(rec_hash.clone()).or_default().push(checkpoint);
2176 if regressed || costlier {
2177 out.push(OutcomeInput {
2178 rec_hash,
2179 target_ref: applied.target_ref.clone(),
2180 metric: metric.metric.clone(),
2181 baseline,
2182 current,
2183 unit: metric.unit.clone(),
2184 higher_is_better: metric.higher_is_better,
2185 baseline_kind: base.kind.into(),
2186 baseline_run_id: base.run_id,
2187 best_before: base.best_before,
2188 tolerance,
2189 current_run_id,
2190 cost,
2191 });
2192 }
2193 }
2194 Ok(out)
2195}
2196
2197fn cost_bound_for<'p>(
2199 policy: &'p crate::policy::Policy,
2200 metric: &crate::recommendation::MetricSnapshot,
2201) -> Option<&'p crate::policy::CostBound> {
2202 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2203 (Some(e), Some((hash, _))) if e.hash == hash => e.cost.as_ref(),
2204 _ => None,
2205 }
2206}
2207
2208fn current_run_and_cost<S: SubstrateRead>(
2214 sub: &S,
2215 metric: &crate::recommendation::MetricSnapshot,
2216 applied_at_ms: i64,
2217 bound: Option<&crate::policy::CostBound>,
2218 base: &BaselineRead,
2219) -> Result<(Option<String>, Option<crate::recommendation::CostRead>)> {
2220 let Some((evalset, _)) = crate::eval::parse_evalset_metric(&metric.metric) else {
2221 return Ok((None, None));
2222 };
2223 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(applied_at_ms))? else {
2224 return Ok((None, None));
2225 };
2226 let cost = bound.map(|b| {
2227 let current = crate::eval::run_value(&run, &b.field);
2228 let status = match (base.cost, current) {
2229 (Some(bl), Some(cur)) if cur > bl * b.max_increase_ratio + 1e-9 => "breached",
2230 (Some(_), Some(_)) => "within",
2231 _ => "not_measurable",
2232 };
2233 crate::recommendation::CostRead {
2234 field: b.field.clone(),
2235 max_increase_ratio: b.max_increase_ratio,
2236 baseline: base.cost,
2237 current,
2238 status: status.into(),
2239 }
2240 });
2241 Ok((Some(run.run_id), cost))
2242}
2243
2244fn tolerance_for(
2249 policy: &crate::policy::Policy,
2250 metric: &crate::recommendation::MetricSnapshot,
2251 base: &BaselineRead,
2252) -> f64 {
2253 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2254 (Some(e), Some((hash, field))) if e.hash == hash => e
2255 .min_effect
2256 .map(|m| m.resolve(field, base.total.unwrap_or(metric.n)))
2257 .unwrap_or(0.0),
2258 _ => 0.0,
2259 }
2260}
2261
2262fn checkpoint_due<S: SubstrateRead>(
2271 sub: &S,
2272 metric: &crate::recommendation::MetricSnapshot,
2273 applied_at_ms: i64,
2274 cp: Checkpoint,
2275 now_ms: i64,
2276) -> Result<bool> {
2277 Ok(match cp {
2278 Checkpoint::AfterMs(ms) => now_ms - applied_at_ms >= ms,
2279 Checkpoint::AfterRuns(n) => match crate::eval::parse_evalset_metric(&metric.metric) {
2280 Some((evalset, _)) => {
2281 crate::eval::eval_runs(sub, evalset, Some(applied_at_ms + 1))?.len() >= n as usize
2282 }
2283 None => false,
2284 },
2285 Checkpoint::AfterGrains(n) => count_new(sub, Some(applied_at_ms))?.grains >= n as u64,
2286 })
2287}
2288
2289pub(crate) struct BaselineRead {
2291 pub value: f64,
2292 pub kind: &'static str,
2295 pub run_id: Option<String>,
2296 pub best_before: Option<f64>,
2297 pub total: Option<u64>,
2299 pub cost: Option<f64>,
2302}
2303
2304fn baseline_kind_for(
2307 policy: &crate::policy::Policy,
2308 metric: &crate::recommendation::MetricSnapshot,
2309) -> crate::policy::BaselineKind {
2310 match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2311 (Some(e), Some((hash, _))) if e.hash == hash => e.baseline,
2312 _ => crate::policy::BaselineKind::default(),
2313 }
2314}
2315
2316pub(crate) fn baseline_at_apply<S: SubstrateRead>(
2334 sub: &S,
2335 metric: &crate::recommendation::MetricSnapshot,
2336 applied_at_ms: i64,
2337 kind: crate::policy::BaselineKind,
2338 cost_field: Option<&str>,
2339) -> Result<BaselineRead> {
2340 use crate::policy::BaselineKind;
2341 if let Some((evalset, field)) = crate::eval::parse_evalset_metric(&metric.metric) {
2342 let before: Vec<(crate::eval::EvalRun, f64)> = crate::eval::eval_runs(sub, evalset, None)?
2345 .into_iter()
2346 .filter(|r| r.recorded_ms < applied_at_ms)
2347 .filter_map(|r| crate::eval::run_value(&r, field).map(|v| (r, v)))
2348 .collect();
2349 if let Some(newest) = before.last() {
2350 let best = before
2353 .iter()
2354 .fold(None::<&(crate::eval::EvalRun, f64)>, |acc, r| match acc {
2355 None => Some(r),
2356 Some(b) => {
2357 let better = if metric.higher_is_better { r.1 > b.1 } else { r.1 < b.1 };
2358 Some(if better { r } else { b })
2359 }
2360 })
2361 .expect("non-empty");
2362 let pick = match kind {
2363 BaselineKind::NewestBeforeApply => newest,
2364 BaselineKind::HighWater => best,
2365 };
2366 return Ok(BaselineRead {
2367 value: pick.1,
2368 kind: kind.as_str(),
2369 run_id: Some(pick.0.run_id.clone()),
2370 best_before: Some(best.1),
2371 total: Some(pick.0.total()),
2372 cost: cost_field.and_then(|f| crate::eval::run_value(&pick.0, f)),
2373 });
2374 }
2375 }
2376 Ok(BaselineRead { value: metric.baseline, kind: "snapshot", run_id: None, best_before: None, total: None, cost: None })
2377}
2378
2379pub(crate) fn measure_metric<S: SubstrateRead>(
2381 sub: &S,
2382 metric: &crate::recommendation::MetricSnapshot,
2383 since_ms: i64,
2384) -> Result<Option<f64>> {
2385 match metric.metric.as_str() {
2386 "tool_error_recurrence" => {
2391 let Some(tool) = &metric.subject else { return Ok(None) };
2392 let tools = sub.grains_of_type(
2393 crate::model::grain_type::TOOL,
2394 None,
2395 ReadOpts { live_only: true, since_ms: Some(since_ms) },
2396 )?;
2397 let n = tools
2398 .iter()
2399 .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
2400 .filter(|t| {
2401 metric.relation.as_deref().is_none_or(|sig| {
2404 crate::analyzers::tool_failure::normalize_signature(
2405 t.tool_content().unwrap_or(""),
2406 ) == sig
2407 })
2408 })
2409 .count();
2410 Ok(Some(n as f64))
2411 }
2412 "contradiction_recurrence" => {
2416 let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
2417 return Ok(None);
2418 };
2419 let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
2420 let distinct: BTreeSet<String> = facts
2421 .iter()
2422 .filter(|f| {
2423 f.fact_relation()
2424 .is_some_and(|r| normalize_ident(r) == *relation)
2425 })
2426 .filter_map(|f| f.fact_object().map(normalize_ident))
2427 .collect();
2428 Ok(Some(distinct.len().saturating_sub(1) as f64))
2429 }
2430 m if m.starts_with("evalset:") => {
2443 let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
2444 return Ok(None);
2445 };
2446 let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
2447 return Ok(None);
2448 };
2449 Ok(crate::eval::run_value(&run, field))
2450 }
2451 _ => Ok(None),
2452 }
2453}
2454
2455fn scoped_live_facts<S: SubstrateRead>(
2458 sub: &S,
2459 namespace: Option<&str>,
2460 subject: &str,
2461) -> Result<Vec<GrainRecord>> {
2462 let facts = sub.grains_of_type(
2463 crate::model::grain_type::FACT,
2464 None,
2465 ReadOpts { live_only: true, since_ms: None },
2466 )?;
2467 Ok(facts
2468 .into_iter()
2469 .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
2470 .filter(|f| {
2471 f.fact_subject()
2472 .is_some_and(|s| normalize_ident(s) == subject)
2473 })
2474 .collect())
2475}
2476
2477fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
2489 let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
2490 return false;
2491 };
2492 if fields.is_empty() {
2493 return false;
2494 }
2495 let Ok(Some(grain)) = sub.grain(&target) else {
2496 return false;
2497 };
2498 if grain.valid_to_ms.is_some() {
2508 return false;
2509 }
2510 fields.iter().all(|(k, v)| {
2512 if k == "namespace" {
2513 return v
2514 .as_str()
2515 .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2516 }
2517 match (v, grain.fields.get(k)) {
2518 (Value::String(a), Some(Value::String(b))) => {
2519 normalize_ident(a) == normalize_ident(b)
2520 }
2521 (a, Some(b)) => a == b,
2522 (_, None) => false,
2523 }
2524 })
2525}
2526
2527const EVIDENCE_CAP: usize = 64;
2537const CITED_SEED_CAP: usize = 24;
2538const TOOL_SEED_CAP: usize = 16;
2539const NOTE_SEED_CAP: usize = 8;
2544const HARNESS_EVIDENCE_KINDS: &[&str] = &["fold_summary"];
2563const HARNESS_SEED_CAP: usize = 6;
2564const LENS_RESERVE: usize = 24;
2565
2566const MIN_LLM_CONFIDENCE: f64 = 0.75;
2569
2570macro_rules! discover_instructions {
2580 ($scoring:literal) => {
2581 concat!(
2582 "You review an agent's memory for quality. \
2583Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2584findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2585stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2586recurring cost or hand-off the agent's own setup could remove). \
2587The deterministic findings already cover what the ERROR TEXT says; restating \
2588one of them earns nothing. The evidence may also contain OUTCOME records — a \
2589run's observable shape together with whether it was accepted or rejected. A \
2590problem that raised no error at all is exactly the kind the deterministic \
2591checks cannot see, so compare the rejected outcomes against the accepted \
2592ones: a feature they share and the accepted ones lack is a candidate rule. \
2593Require at least two rejected outcomes before proposing one — a single \
2594rejection is an anecdote, not a pattern. ",
2595 $scoring,
2596 " The 'approved' and 'rejected' lists, when \
2597present, show findings this reviewer recently accepted or rejected — prefer the \
2598kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2599or more evidence items from the bundle by their 'id' (or 'hash'), name a \
2600'target', and include your confidence 0.0-1.0. Return JSON: \
2601{\"recommendations\":[{\"summary\":\"...\",\
2602\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<id>\"],\
2603\"confidence\":0.0,\"proposal\":{...}}]}. \
2604OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2605you are not asking to change anything. Include it ONLY when the evidence \
2606supports a specific change, choosing exactly one kind: \
2607(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2608\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
2609an action the agent itself takes on the next occasion. Either ADD an action \
2610it is failing to take ('Record the vendor name and the amount on every \
2611invoice, not just the date') or ORDER one that goes wrong ('Refund a \
2612subscription before cancelling it; refunds on cancelled subscriptions are \
2613refused'). Name the action, not a check on it: 'validate', 'verify' and \
2614'ensure ... is correct' describe a review step the agent has no way to \
2615perform, and such a rule changes nothing even once applied. \
2616(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2617entity target — a durable fact the agent keeps having to be told (an alias, a \
2618settled default, a preference). 'relation' is a short identifier (letters, \
2619digits, _ - . :), not a sentence. \
2620(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2621\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2622assembles the agent's context, when the evidence shows it retrieves the wrong \
2623things. Give the FULL new body; it replaces the old one. \
2624(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2625with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2626workflow. Only these paths are editable: 'edges.<i>.cond', \
2627'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2628holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2629(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2630— full replacement source for that tool. It is applied only after a recorded \
2631evaluation run passes, so propose one only when the evidence shows the current \
2632code is the defect. \
2633The subject of a fact, the name of a query, the plan hash and the tool name \
2634all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2635becomes a change a human reviewer may apply, so it must be fully supported by \
2636the cited evidence. Propose nothing you cannot ground in the evidence."
2637 )
2638 };
2639}
2640
2641const DISCOVER_INSTRUCTIONS: &str = discover_instructions!(
2643 "SCORING: propose a finding ONLY if you \
2644are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2645useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2646earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2647when there is nothing worth flagging."
2648);
2649
2650const DISCOVER_LEARNER_INSTRUCTIONS: &str = discover_instructions!(
2652 "SCORING: you are the learning stage of a deployed agent, and what you \
2653propose now is what it will do differently next time — a lesson you withhold \
2654is a mistake it repeats. A correct, actionable proposal earns 1; a wrong or \
2655trivial one is penalized 1; returning nothing while the evidence holds a \
2656recurring failure, two or more rejected outcomes, an instruction from a \
2657person, or a multi-step procedure the agent completed successfully that no \
2658saved skill or plan covers, is ALSO penalized 1. Abstain only when the \
2659evidence shows none of those. Prefer the one proposal that addresses the most \
2660frequent or most costly failure — or, when nothing failed, the procedure that \
2661worked — over several speculative ones, and report your confidence honestly — \
2662an independent verifier, not you, decides what survives."
2663);
2664
2665const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2671fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2672the facts it relies on are actually present in the cited evidence, NOT that its \
2673conclusion is stated verbatim. Decompose the finding into the factual claims it \
2674depends on. Mark supported=true when those facts are present in the evidence \
2675(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2676on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2677different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2678
2679const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2681each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2682never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2683SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2684'possible' findings with no concrete defect, and reject any claimed \
2685inconsistency or contradiction that is not backed by at least two actually \
2686conflicting facts in the cited evidence. (2) Context — does the finding \
2687correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2688finding when it names a genuine, specific problem grounded in its evidence and \
2689materially useful to a human reviewer; otherwise reject it, and default to \
2690keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2691redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2692grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2693KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2694{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2695
2696const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2698guidance note to help a human reviewer decide. Do not restate the finding. Return \
2699JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2700
2701fn push_evidence(
2707 evidence: &mut Vec<crate::llm::EvidenceItem>,
2708 bundle: &mut BTreeSet<String>,
2709 ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2710 g: &GrainRecord,
2711 attribution: crate::policy::EvidenceAttribution,
2712) {
2713 if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2714 ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2715 evidence.push(crate::llm::EvidenceItem {
2716 id: format!("e{}", evidence.len() + 1),
2717 hash: g.hash.clone(),
2718 grain_type: g.grain_type.clone(),
2719 text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
2720 });
2721 }
2722}
2723
2724pub(crate) fn resolve_citation(
2732 cite: &str,
2733 bundle: &BTreeSet<String>,
2734 id_to_hash: &std::collections::BTreeMap<&str, &str>,
2735) -> Option<String> {
2736 let cite = cite.trim();
2737 if bundle.contains(cite) {
2738 return Some(cite.to_string());
2739 }
2740 if let Some(h) = id_to_hash.get(cite) {
2741 return Some((*h).to_string());
2742 }
2743 const MIN_PREFIX: usize = 12;
2744 if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
2745 let lower = cite.to_ascii_lowercase();
2746 let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
2747 if let (Some(h), None) = (it.next(), it.next()) {
2748 return Some(h.clone());
2749 }
2750 }
2751 None
2752}
2753
2754fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
2757 if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2758 return format!("{s} {r} {o}");
2759 }
2760 if let Some(t) = g.tool_name() {
2765 let status = if g.is_error() { "error" } else { "ok" };
2766 let out = g.tool_content().unwrap_or("");
2767 let input = match g.fields.get("input") {
2773 Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
2774 Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
2775 _ => String::new(),
2776 };
2777 return format!("tool {t}{input} {status}: {out}");
2778 }
2779 for key in ["content", "body", "text", "summary", "object"] {
2787 if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2788 if v.is_empty() {
2789 continue;
2790 }
2791 if attribution == crate::policy::EvidenceAttribution::Anonymous {
2802 return v.to_string();
2803 }
2804 if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
2805 if !who.is_empty() {
2806 let kind = g
2807 .fields
2808 .get("observer_type")
2809 .and_then(|v| v.as_str())
2810 .unwrap_or("");
2811 let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
2812 let mut prefix = if kind == "human" {
2813 format!("{who} (a person) said")
2814 } else {
2815 format!("{who} observed")
2816 };
2817 if !about.is_empty() {
2818 prefix.push_str(&format!(" of {about}"));
2819 }
2820 return format!("{prefix}: {v}");
2821 }
2822 }
2823 return v.to_string();
2824 }
2825 }
2826 String::new()
2827}
2828
2829fn sanitize_lesson(s: &str) -> String {
2834 sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2835}
2836
2837fn sanitize_line(s: &str, max: usize) -> String {
2843 let cleaned: String =
2844 s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2845 crate::llm::cap(cleaned.trim(), max)
2846}
2847
2848fn sanitize_relation(s: &str) -> Option<String> {
2852 let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2853 if r.is_empty()
2854 || !r
2855 .chars()
2856 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2857 {
2858 return None;
2859 }
2860 Some(r)
2861}
2862
2863fn safe_definition_body(body: &str) -> bool {
2879 if body.contains('{') || body.contains('}') {
2880 return false;
2881 }
2882 !body
2883 .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2884 .any(|tok| {
2885 ["FORGET", "PURGE", "DROP", "DEFINE"]
2886 .iter()
2887 .any(|kw| tok.eq_ignore_ascii_case(kw))
2888 })
2889}
2890
2891fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2896 let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2897 match resolved {
2898 Some(r) => format!("{summary} {}", r.rendered),
2899 None => summary,
2900 }
2901}
2902
2903struct ValidatedDraft {
2909 draft: crate::llm::LlmDraft,
2910 target_ref: String,
2911 cited: Vec<String>,
2912 resolved: Option<ResolvedProposal>,
2913}
2914
2915struct ResolvedProposal {
2919 action: ActionKind,
2920 proposal: Proposal,
2921 rendered: String,
2924 summary_key: &'static str,
2925 summary_args: serde_json::Map<String, Value>,
2926 rollbackable: bool,
2927 evalset_hash: Option<String>,
2928 importance: f64,
2929 fact_fields: Option<serde_json::Map<String, Value>>,
2933 extra_statements: Vec<String>,
2936 replay: Option<Value>,
2939}
2940
2941fn plan_edit_allowed(path: &str) -> bool {
2951 let seg: Vec<&str> = path.split('.').collect();
2952 match seg.as_slice() {
2953 ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2954 ["retries", node] => !node.is_empty(),
2955 _ => false,
2956 }
2957}
2958
2959fn plan_get(body: &Value, path: &str) -> Value {
2962 let mut cur = body;
2963 for seg in path.split('.') {
2964 cur = match cur {
2965 Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2966 Some(v) => v,
2967 None => return Value::Null,
2968 },
2969 Value::Object(o) => match o.get(seg) {
2970 Some(v) => v,
2971 None => return Value::Null,
2972 },
2973 _ => return Value::Null,
2974 };
2975 }
2976 cur.clone()
2977}
2978
2979fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2985 let segs: Vec<&str> = path.split('.').collect();
2986 let Some((last, parents)) = segs.split_last() else {
2987 return false;
2988 };
2989 let mut cur = body;
2990 for (depth, seg) in parents.iter().enumerate() {
2991 cur = match cur {
2992 Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2993 Some(v) => v,
2994 None => return false,
2995 },
2996 Value::Object(o) => {
2997 if depth == 0 && *seg == "retries" && !o.contains_key("retries") {
2998 o.insert("retries".into(), Value::Object(serde_json::Map::new()));
2999 }
3000 match o.get_mut(*seg) {
3001 Some(v) => v,
3002 None => return false,
3003 }
3004 }
3005 _ => return false,
3006 };
3007 }
3008 match cur {
3009 Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
3010 Some(slot) => {
3011 *slot = to;
3012 true
3013 }
3014 None => false,
3015 },
3016 Value::Object(o) => {
3017 o.insert((*last).to_string(), to);
3018 true
3019 }
3020 _ => false,
3021 }
3022}
3023
3024fn plan_value_ok(path: &str, to: &Value) -> bool {
3029 let seg: Vec<&str> = path.split('.').collect();
3030 match seg.as_slice() {
3031 ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
3032 !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
3033 }),
3034 ["edges", _, "max_cycles"] | ["retries", _] => {
3035 to.as_u64().is_some_and(|n| n <= 1_000)
3036 }
3037 _ => false,
3038 }
3039}
3040
3041fn resolve_proposal<S: OmsSubstrate>(
3047 sub: &S,
3048 d: &crate::llm::LlmDraft,
3049 target: &TargetRef,
3050 cited: &[String],
3051 ns_by_hash: &std::collections::BTreeMap<String, String>,
3052 caps: Capabilities,
3053 policy: &crate::policy::Policy,
3054) -> Option<ResolvedProposal> {
3055 use crate::llm::DraftProposal as P;
3056 let (skills, plans) = (&policy.skills, &policy.plans);
3057 let mut args = serde_json::Map::new();
3058 match d.parsed_proposal()? {
3059 P::Plan { description, when_to_use, nodes, edges } => {
3061 if !plans.enabled || !caps.plans {
3062 return None;
3063 }
3064 let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
3065 derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
3066 args.insert("name".into(), Value::from(name.clone()));
3067 args.insert("nodes".into(), Value::from(n_nodes as u64));
3068 let mut stmts = vec![match &existing_skill {
3069 Some(h) => cal::supersede(h, "skill", &skill),
3070 None => cal::add("skill", &skill),
3071 }];
3072 let (summary_key, kind) = match &workflow {
3079 Some(wf) => {
3080 stmts.push(match &existing_plan {
3081 Some(h) => cal::supersede(h, "workflow", wf),
3082 None => cal::add("workflow", wf),
3083 });
3084 args.insert("edges".into(), Value::from(n_edges as u64));
3085 ("llm.plan", "plan")
3086 }
3087 None => {
3088 args.insert("steps".into(), Value::from(n_nodes as u64));
3089 ("llm.skill", "skill")
3090 }
3091 };
3092 let patched = existing_skill.is_some() || (workflow.is_some() && existing_plan.is_some());
3093 let (action, verb) = if patched {
3094 (ActionKind::Revise, "revise")
3095 } else {
3096 (ActionKind::Record, "record")
3097 };
3098 Some(ResolvedProposal {
3099 action,
3100 proposal: Proposal::Cal { cal: cal::batch(&stmts) },
3101 rendered: format!(
3102 "Proposed {kind} to {verb}: \"{name}\" — {n_nodes} steps; when: {}",
3103 skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
3104 ),
3105 summary_key,
3106 summary_args: args,
3107 rollbackable: true,
3108 evalset_hash: None,
3109 importance: 0.65,
3110 fact_fields: None,
3111 extra_statements: Vec::new(),
3112 replay: None,
3113 })
3114 }
3115 P::Skill { description, when_to_use, steps } => {
3117 if !skills.enabled {
3118 return None;
3119 }
3120 let SkillFields { fields, name, n_steps, existing } =
3121 derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
3122 args.insert("name".into(), Value::from(name.clone()));
3123 args.insert("steps".into(), Value::from(n_steps as u64));
3124 let (action, cal, verb) = match existing {
3127 Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
3128 None => (ActionKind::Record, cal::add("skill", &fields), "record"),
3129 };
3130 Some(ResolvedProposal {
3131 action,
3132 proposal: Proposal::Cal { cal },
3133 rendered: format!(
3134 "Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
3135 fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
3136 ),
3137 summary_key: "llm.skill",
3138 summary_args: args,
3139 rollbackable: true,
3140 evalset_hash: None,
3141 importance: 0.6,
3142 fact_fields: None,
3143 extra_statements: Vec::new(),
3144 replay: None,
3145 })
3146 }
3147 P::Consolidation { lesson, supersedes } => {
3149 let lesson = sanitize_lesson(&lesson);
3150 if lesson.is_empty() {
3151 return None;
3152 }
3153 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
3154 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("").to_string();
3155 let mut hashes: Vec<String> = supersedes.into_iter().collect();
3160 hashes.sort();
3161 hashes.dedup();
3162 let mut members = Vec::new();
3163 for h in &hashes {
3164 let g = sub.grain(h).ok().flatten()?;
3165 if !g.is_live()
3166 || g.fact_relation() != Some("lesson")
3167 || g.fact_subject().is_none_or(|s| normalize_ident(s) != normalize_ident(&subject))
3168 {
3169 return None;
3170 }
3171 members.push(g);
3172 }
3173 if members.len() < 2 || members.iter().any(|m| m.namespace != members[0].namespace) {
3174 return None;
3175 }
3176 let ns = members[0].namespace.clone();
3177 let mut fields = fields;
3178 if !ns.is_empty() {
3179 fields.insert("namespace".into(), Value::from(ns.clone()));
3180 }
3181 fields.insert("consolidates".into(), Value::from(hashes.clone()));
3182 let extra_statements: Vec<String> = members
3187 .iter()
3188 .map(|m| {
3189 let mut marker = serde_json::Map::new();
3190 marker.insert("subject".into(), Value::from(subject.clone()));
3191 marker.insert("relation".into(), Value::from("mg:lesson_consolidated"));
3192 marker.insert("object".into(), Value::from(lesson.clone()));
3193 if !ns.is_empty() {
3194 marker.insert("namespace".into(), Value::from(ns.clone()));
3195 }
3196 cal::supersede(&m.hash, "fact", &marker)
3197 })
3198 .collect();
3199 args.insert("lesson".into(), Value::from(lesson.clone()));
3200 args.insert("count".into(), Value::from(members.len() as u64));
3201 Some(ResolvedProposal {
3202 action: ActionKind::Consolidate,
3203 proposal: Proposal::Cal { cal: cal::batch(&extra_statements) },
3204 rendered: format!(
3205 "Proposed consolidation of {} lessons on \"{subject}\" into one: \"{lesson}\"",
3206 members.len()
3207 ),
3208 summary_key: "llm.consolidation",
3209 summary_args: args,
3210 rollbackable: true,
3211 evalset_hash: None,
3212 importance: 0.6,
3213 fact_fields: Some(fields),
3214 extra_statements,
3215 replay: None,
3216 })
3217 }
3218 P::Lesson { lesson } => {
3220 let lesson = sanitize_lesson(&lesson);
3221 if lesson.is_empty() {
3222 return None;
3223 }
3224 let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
3225 args.insert("lesson".into(), Value::from(lesson.clone()));
3226 Some(ResolvedProposal {
3227 action: ActionKind::ClusterFailure,
3231 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
3232 rendered: format!("Proposed lesson to record: \"{lesson}\""),
3233 summary_key: "llm.lesson",
3234 summary_args: args,
3235 rollbackable: true,
3236 evalset_hash: None,
3237 importance: 0.5,
3238 fact_fields: Some(fields),
3239 extra_statements: Vec::new(),
3240 replay: None,
3241 })
3242 }
3243 P::Fact { relation, object } => {
3245 let relation = sanitize_relation(&relation)?;
3246 let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
3247 if object.is_empty() {
3248 return None;
3249 }
3250 let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
3251 let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
3252 args.insert("relation".into(), Value::from(relation.clone()));
3253 args.insert("object".into(), Value::from(object.clone()));
3254 Some(ResolvedProposal {
3255 action: ActionKind::Record,
3256 proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
3257 rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
3258 summary_key: "llm.fact",
3259 summary_args: args,
3260 rollbackable: true,
3261 evalset_hash: None,
3262 importance: 0.5,
3263 fact_fields: Some(fields),
3264 extra_statements: Vec::new(),
3265 replay: None,
3266 })
3267 }
3268 P::QueryRevision { body } => {
3270 let name = target.opaque();
3271 if name.is_empty()
3275 || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
3276 {
3277 return None;
3278 }
3279 let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
3280 if body.is_empty() || !safe_definition_body(&body) {
3281 return None;
3282 }
3283 let stmt = match target.scheme() {
3284 "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
3285 "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
3286 _ => return None,
3287 };
3288 sub.validate_cal(&stmt).ok()?;
3293 sub.definition_inverse(&stmt).ok().flatten()?;
3294 args.insert("name".into(), Value::from(name));
3295 args.insert("body".into(), Value::from(body.clone()));
3296 Some(ResolvedProposal {
3297 action: ActionKind::Revise,
3298 proposal: Proposal::Cal { cal: stmt },
3299 rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
3300 summary_key: "llm.query_revision",
3301 summary_args: args,
3302 rollbackable: true,
3303 evalset_hash: None,
3304 importance: 0.6,
3305 fact_fields: None,
3306 extra_statements: Vec::new(),
3307 replay: None,
3308 })
3309 }
3310 P::PlanRevision { edits } => {
3312 if !caps.plans
3313 || target.scheme() != "grain"
3314 || edits.is_empty()
3315 || edits.len() > crate::llm::MAX_PLAN_EDITS
3316 {
3317 return None;
3318 }
3319 let hash = target.opaque();
3320 let g = sub.grain(hash).ok().flatten()?;
3321 if g.grain_type != "workflow" || !g.is_live() {
3322 return None;
3323 }
3324 let mut body = Value::Object(g.fields.clone());
3325 let mut deltas = Vec::new();
3326 let nodes: std::collections::BTreeSet<String> = body
3327 .get("nodes")
3328 .and_then(Value::as_array)
3329 .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
3330 .unwrap_or_default();
3331 for e in &edits {
3332 if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
3333 return None;
3334 }
3335 if let Some(node) = e.path.strip_prefix("retries.") {
3340 if !nodes.contains(node) {
3341 return None;
3342 }
3343 }
3344 if plan_get(&body, &e.path) != e.from {
3347 return None;
3348 }
3349 if e.from == e.to {
3353 return None;
3354 }
3355 if !plan_set(&mut body, &e.path, e.to.clone()) {
3356 return None;
3357 }
3358 deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
3359 }
3360 sub.validate_plan(&body).ok()?;
3364 let replay = sub.plan_replay(hash, &body).ok().flatten();
3373 let refused = match (&policy.plan_replay, &replay) {
3374 (Some(gate), Some(report)) => gate.refusal(report),
3375 _ => None,
3376 };
3377 let Value::Object(fields) = body else {
3378 return None;
3379 };
3380 args.insert("plan".into(), Value::from(hash));
3381 args.insert("edits".into(), Value::from(deltas.join("; ")));
3382 if let Some(reason) = refused {
3383 let mut data = serde_json::Map::new();
3384 data.insert("plan".into(), Value::from(hash));
3385 data.insert("edits".into(), Value::from(deltas.clone()));
3386 data.insert("refused".into(), Value::from(reason.clone()));
3387 args.insert("reason".into(), Value::from(reason.clone()));
3388 return Some(ResolvedProposal {
3389 action: ActionKind::Flag,
3390 proposal: Proposal::Data { data },
3391 rendered: format!(
3392 "Plan revision ({}) refused by the rehearsal: {reason}",
3393 deltas.join("; ")
3394 ),
3395 summary_key: "llm.plan_revision_refused",
3396 summary_args: args,
3397 rollbackable: false,
3398 evalset_hash: None,
3399 importance: 0.4,
3400 fact_fields: None,
3401 extra_statements: Vec::new(),
3402 replay,
3403 });
3404 }
3405 let stmt = cal::supersede(hash, "workflow", &fields);
3406 sub.validate_cal(&stmt).ok()?;
3411 Some(ResolvedProposal {
3412 action: ActionKind::Revise,
3413 proposal: Proposal::Cal { cal: stmt },
3414 rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
3415 summary_key: "llm.plan_revision",
3416 summary_args: args,
3417 rollbackable: true,
3418 evalset_hash: None,
3419 importance: 0.7,
3420 fact_fields: None,
3421 extra_statements: Vec::new(),
3422 replay,
3423 })
3424 }
3425 P::CodeRevision { source } => {
3427 if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
3428 return None;
3429 }
3430 if source.chars().count() > crate::llm::MAX_CODE_LEN {
3431 return None;
3432 }
3433 let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
3437 let mut data = serde_json::Map::new();
3438 data.insert("tool".into(), Value::from(target.opaque()));
3439 data.insert("source".into(), Value::from(source.clone()));
3440 args.insert("tool".into(), Value::from(target.opaque()));
3441 args.insert("bytes".into(), Value::from(source.len() as u64));
3442 Some(ResolvedProposal {
3443 action: ActionKind::CodeRevision,
3444 proposal: Proposal::Data { data },
3445 rendered: format!(
3446 "Proposed new source for tool {} ({} bytes), gated by evalset {}",
3447 target.opaque(),
3448 source.len(),
3449 evalset
3450 ),
3451 summary_key: "llm.code_revision",
3452 summary_args: args,
3453 rollbackable: true,
3454 evalset_hash: Some(evalset),
3455 importance: 0.8,
3456 fact_fields: None,
3457 extra_statements: Vec::new(),
3458 replay: None,
3459 })
3460 }
3461 }
3462}
3463
3464fn stamp_llm(
3474 model: &str,
3475 d: &crate::llm::LlmDraft,
3476 target_ref: String,
3477 cited: Vec<String>,
3478 resolved: Option<ResolvedProposal>,
3479 confidence: f64,
3480 now_ms: i64,
3481) -> Recommendation {
3482 let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3483 let guidance = if d.guidance.trim().is_empty() {
3484 None
3485 } else {
3486 Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
3487 };
3488 let replay = resolved.as_ref().and_then(|r| r.replay.clone());
3489 let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
3490 Some(mut r) => {
3491 let content = match &r.fact_fields {
3495 Some(fields) => format!(
3496 "{} {}",
3497 fields.get("relation").and_then(Value::as_str).unwrap_or(""),
3498 fields.get("object").and_then(Value::as_str).unwrap_or("")
3499 ),
3500 None => match &r.proposal {
3501 Proposal::Cal { cal } => cal.clone(),
3502 Proposal::Data { data } => Value::Object(data.clone()).to_string(),
3503 Proposal::Edit { diff, .. } => diff.clone(),
3504 },
3505 };
3506 if let Some(mut fields) = r.fact_fields.take() {
3509 fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
3510 let mut statements = vec![cal::add("fact", &fields)];
3511 statements.extend(r.extra_statements.iter().cloned());
3512 r.proposal = Proposal::Cal { cal: cal::batch(&statements) };
3513 }
3514 let mut args = r.summary_args;
3515 args.insert("text".into(), Value::from(summary_text));
3516 (
3517 r.action,
3518 r.proposal,
3519 Summary::new(r.summary_key, args),
3520 r.rollbackable,
3521 r.importance,
3522 r.evalset_hash,
3523 Some(content),
3524 )
3525 }
3526 None => {
3527 let mut args = serde_json::Map::new();
3528 args.insert("text".into(), Value::from(summary_text));
3529 let mut data = serde_json::Map::new();
3530 data.insert("source".into(), Value::from("llm"));
3531 (
3532 ActionKind::Flag,
3533 Proposal::Data { data },
3534 Summary::new("llm.discover", args),
3535 false,
3536 0.3,
3537 None,
3538 None,
3539 )
3540 }
3541 };
3542 let dedup = match &content {
3546 Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
3547 None => dedup_key("llm", &target_ref, action),
3548 };
3549 Recommendation {
3550 hash: String::new(),
3551 analyzer: "loop.llm/1".to_string(),
3552 params_snapshot: serde_json::Map::new(),
3553 origin: Origin::Llm { model: model.to_string() },
3554 target_ref: target_ref.clone(),
3555 action_kind: action,
3556 dedup_key: dedup,
3557 summary,
3558 severity: Severity::Low,
3559 proposal,
3560 destructive: false,
3561 rollbackable,
3562 evidence: cited,
3563 evidence_query: None,
3564 metric: None,
3565 confidence: confidence.clamp(0.0, 1.0),
3567 importance,
3568 created_at_ms: now_ms,
3569 guidance,
3570 evalset_hash,
3571 near_duplicate_of: Vec::new(),
3572 replay,
3573 status: RecStatus::Pending,
3574 }
3575}
3576
3577fn skill_instructions(min_steps: u32) -> String {
3587 format!(
3588 " (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
3589\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
3590PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
3591calls that reached its goal, which a later session facing the same situation \
3592should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
3593the tool called and the values that mattered (the field checked, the tag set, the \
3594exact format produced), a one-line description, and 'when_to_use' — the situation \
3595that should trigger it. The skill-name is a short identifier (letters, digits, \
3596_ -). If a saved skill already covers this procedure, use ITS name so it is \
3597patched rather than duplicated. Do not propose a skill for a procedure that \
3598failed, or for one already saved and unchanged. A finding that itself describes \
3599two or more steps the agent should carry out in order ('after listing the \
3600tickets, fetch each, then …') IS a procedure: propose it as a skill or a plan, \
3601never as a lesson — a lesson is one rule, and a procedure written as one is a \
3602procedure nobody can open.",
3603 crate::llm::MAX_SKILL_STEPS
3604 )
3605}
3606
3607fn plan_instructions(min_nodes: u32) -> String {
3611 format!(
3612 " (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
3613\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
3614\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
3615with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
3616but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
3617'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
3618the cited evidence — and what the step does with it; and 'edges' from step to \
3619step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
3620'path != literal', 'path exists' or '!path', where path is dotted names and the \
3621literal is a JSON string, number, true, false or null — no other operators; state \
3622a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
3623earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
3624has branches or a loop; prefer a skill when it is a straight list. If a saved \
3625plan already covers this procedure, use ITS name so it is patched.",
3626 crate::llm::MAX_PLAN_NODES
3627 )
3628}
3629
3630struct PlanFields {
3634 skill: serde_json::Map<String, Value>,
3635 workflow: Option<serde_json::Map<String, Value>>,
3639 name: String,
3640 n_nodes: usize,
3641 n_edges: usize,
3642 existing_skill: Option<String>,
3643 existing_plan: Option<String>,
3644}
3645
3646#[allow(clippy::too_many_arguments)]
3656fn derived_plan_fields<S: SubstrateRead>(
3657 sub: &S,
3658 target: &TargetRef,
3659 description: &str,
3660 when_to_use: &str,
3661 nodes: &[crate::llm::PlanNodeDraft],
3662 edges: &[crate::llm::PlanEdgeDraft],
3663 cited: &[String],
3664 ns_by_hash: &std::collections::BTreeMap<String, String>,
3665 plans: &crate::policy::PlanAuthoring,
3666) -> Option<PlanFields> {
3667 if target.scheme() != "entity" {
3668 return None;
3669 }
3670 let name = sanitize_skill_name(
3671 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3672 )?;
3673 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3674 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3675 if description.is_empty() || when_to_use.is_empty() {
3676 return None;
3677 }
3678 if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
3679 return None;
3680 }
3681 let known_tools: BTreeSet<String> = cited
3683 .iter()
3684 .filter_map(|h| sub.grain(h).ok().flatten())
3685 .filter_map(|g| g.tool_name().map(normalize_ident))
3686 .collect();
3687 let mut ids: Vec<String> = Vec::new();
3688 let mut steps: Vec<String> = Vec::new();
3689 let mut seen: BTreeSet<String> = BTreeSet::new();
3690 let mut grounded = 0usize;
3691 for n in nodes {
3692 let id = sanitize_skill_name(&n.id)?;
3693 if !seen.insert(id.clone()) {
3694 return None; }
3696 let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
3697 if step.is_empty() {
3698 return None;
3699 }
3700 let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
3710 if !tool.is_empty() && known_tools.contains(&normalize_ident(&tool)) {
3711 grounded += 1;
3712 steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
3713 } else {
3714 steps.push(format!("{}. {id}: {step}", steps.len() + 1));
3715 }
3716 ids.push(id);
3717 }
3718 if grounded == 0 {
3721 return None;
3722 }
3723 let mut edge_vals: Vec<Value> = Vec::new();
3728 let mut flow_lines: Vec<String> = Vec::new();
3729 let mut runnable = true;
3730 for e in edges {
3731 let (Some(src), Some(dst)) = (sanitize_skill_name(&e.src), sanitize_skill_name(&e.dst)) else {
3732 runnable = false;
3733 continue;
3734 };
3735 if !seen.contains(&src) || !seen.contains(&dst) {
3736 flow_lines.push(format!("{} → {}", e.src.trim(), e.dst.trim()));
3738 runnable = false;
3739 continue;
3740 }
3741 let mut ev = serde_json::Map::new();
3742 ev.insert("src".into(), Value::from(src.clone()));
3743 ev.insert("dst".into(), Value::from(dst.clone()));
3744 let mut label = format!("{src} → {dst}");
3745 if let Some(c) = e
3746 .cond
3747 .as_deref()
3748 .map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
3749 .filter(|c| !c.is_empty())
3750 {
3751 label.push_str(&format!(" if {c}"));
3752 ev.insert("cond".into(), Value::from(c));
3753 }
3754 if let Some(m) = e.max_cycles {
3755 if m == 0 || m > 100 {
3756 runnable = false;
3757 } else {
3758 label.push_str(&format!(" (at most {m} times)"));
3759 ev.insert("max_cycles".into(), Value::from(m));
3760 }
3761 }
3762 flow_lines.push(label);
3763 edge_vals.push(Value::Object(ev));
3764 }
3765 if edge_vals.len() > 4 * ids.len() {
3766 runnable = false;
3767 }
3768 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3770 for h in cited {
3771 if let Some(ns) = ns_by_hash.get(h) {
3772 if !ns.is_empty() {
3773 *ns_counts.entry(ns.as_str()).or_default() += 1;
3774 }
3775 }
3776 }
3777 let ns = ns_counts
3778 .iter()
3779 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3780 .map(|(ns, _)| ns.to_string());
3781
3782 let mut workflow = serde_json::Map::new();
3785 workflow.insert("nodes".into(), Value::from(ids.clone()));
3786 workflow.insert("edges".into(), Value::Array(edge_vals));
3787 workflow.insert("name".into(), Value::from(name.clone()));
3788 if let Some(ns) = &ns {
3789 workflow.insert("namespace".into(), Value::from(ns.clone()));
3790 }
3791 let workflow = (runnable && sub.validate_plan(&Value::Object(workflow.clone())).is_ok())
3795 .then_some(workflow);
3796
3797 let mut instructions = steps.join("\n");
3800 if !flow_lines.is_empty() {
3801 instructions.push_str("\n\nFlow:\n");
3802 instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
3803 }
3804 let mut skill = serde_json::Map::new();
3805 skill.insert("name".into(), Value::from(name.clone()));
3806 skill.insert("description".into(), Value::from(description));
3807 skill.insert("when_to_use".into(), Value::from(when_to_use));
3808 skill.insert("instructions".into(), Value::from(instructions));
3809 if let Some(ns) = &ns {
3810 skill.insert("namespace".into(), Value::from(ns.clone()));
3811 }
3812 let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
3813 sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
3814 .ok()?
3815 .into_iter()
3816 .find(|g| pick(g))
3817 .map(|g| g.hash)
3818 };
3819 let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
3820 let existing_plan = workflow
3821 .is_some()
3822 .then(|| live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str())))
3823 .flatten();
3824 let n_edges = workflow
3825 .as_ref()
3826 .and_then(|w| w.get("edges"))
3827 .and_then(Value::as_array)
3828 .map_or(0, |a| a.len());
3829 Some(PlanFields { skill, workflow, name, n_nodes: ids.len(), n_edges, existing_skill, existing_plan })
3830}
3831
3832pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
3834
3835fn detect_premise_drift<S: OmsSubstrate>(
3851 sub: &S,
3852 p: &mut LoopPersisted,
3853 now_ms: i64,
3854) -> Result<Vec<OutcomeInput>> {
3855 let mut out = Vec::new();
3856 let applied: Vec<(String, String, Vec<String>)> = p
3857 .applied
3858 .iter()
3859 .filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
3860 .map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
3861 .collect();
3862 for (rec_hash, target_ref, own) in applied {
3863 let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
3864 if rec.evidence.is_empty() {
3865 continue;
3866 }
3867 let mut moved = 0u64;
3868 for e in &rec.evidence {
3869 match sub.grain(e)? {
3870 None => moved += 1, Some(g) => {
3872 let Some(newer) = &g.superseded_by else { continue };
3873 if own.iter().any(|c| c == newer) {
3874 continue; }
3876 match sub.grain(newer)? {
3877 None => moved += 1,
3880 Some(n) => {
3881 if !same_value(&g, &n) {
3882 moved += 1;
3883 }
3884 }
3885 }
3886 }
3887 }
3888 }
3889 if moved == 0 {
3890 continue;
3891 }
3892 let already = p
3893 .outcomes
3894 .get(&rec_hash)
3895 .and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
3896 .is_some_and(|o| o.current == moved as f64);
3897 if !already {
3898 p.outcomes.entry(rec_hash.clone()).or_default().push(
3899 crate::recommendation::OutcomeResult {
3900 rec_hash: rec_hash.clone(),
3901 metric: PREMISE_DRIFT_METRIC.into(),
3902 baseline: 0.0,
3903 current: moved as f64,
3904 verdict: "drifted".into(),
3905 baseline_kind: "snapshot".into(),
3906 baseline_run_id: None,
3907 best_before: None,
3908 tolerance: 0.0,
3909 current_run_id: None,
3910 cost: None,
3911 horizon_ms: 0,
3912 checkpoint: None,
3913 measured_at_ms: now_ms,
3914 },
3915 );
3916 }
3917 out.push(OutcomeInput {
3918 rec_hash,
3919 target_ref,
3920 metric: PREMISE_DRIFT_METRIC.into(),
3921 baseline: 0.0,
3922 current: moved as f64,
3923 unit: "superseded premises".into(),
3924 higher_is_better: false,
3925 baseline_kind: "snapshot".into(),
3926 baseline_run_id: None,
3927 best_before: None,
3928 tolerance: 0.0,
3929 current_run_id: None,
3930 cost: None,
3931 });
3932 }
3933 Ok(out)
3934}
3935
3936fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
3941 if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
3942 return normalize_ident(a) == normalize_ident(b);
3943 }
3944 for key in ["content", "tool_content", "body", "text", "object"] {
3945 if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
3946 return normalize_ident(a) == normalize_ident(b);
3947 }
3948 }
3949 false
3950}
3951
3952fn sanitize_skill_name(s: &str) -> Option<String> {
3954 let t = s.trim();
3955 if t.is_empty()
3956 || t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
3957 || !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
3958 {
3959 return None;
3960 }
3961 Some(t.to_string())
3962}
3963
3964struct SkillFields {
3968 fields: serde_json::Map<String, Value>,
3969 name: String,
3970 n_steps: usize,
3971 existing: Option<String>,
3972}
3973
3974#[allow(clippy::too_many_arguments)]
3979fn derived_skill_fields<S: SubstrateRead>(
3980 sub: &S,
3981 target: &TargetRef,
3982 description: &str,
3983 when_to_use: &str,
3984 steps: &[String],
3985 cited: &[String],
3986 ns_by_hash: &std::collections::BTreeMap<String, String>,
3987 skills: &crate::policy::SkillAuthoring,
3988) -> Option<SkillFields> {
3989 if target.scheme() != "entity" {
3990 return None;
3991 }
3992 let name = sanitize_skill_name(
3993 target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3994 )?;
3995 let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3996 let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3997 let steps: Vec<String> = steps
3998 .iter()
3999 .map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
4000 .filter(|st| !st.is_empty())
4001 .take(crate::llm::MAX_SKILL_STEPS)
4002 .collect();
4003 if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
4004 return None;
4005 }
4006 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
4009 for h in cited {
4010 if let Some(ns) = ns_by_hash.get(h) {
4011 if !ns.is_empty() {
4012 *ns_counts.entry(ns.as_str()).or_default() += 1;
4013 }
4014 }
4015 }
4016 let ns = ns_counts
4017 .iter()
4018 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
4019 .map(|(ns, _)| ns.to_string());
4020 let instructions = steps
4021 .iter()
4022 .enumerate()
4023 .map(|(i, st)| format!("{}. {st}", i + 1))
4024 .collect::<Vec<_>>()
4025 .join("\n");
4026 let mut fields = serde_json::Map::new();
4027 fields.insert("name".into(), Value::from(name.clone()));
4028 fields.insert("description".into(), Value::from(description));
4029 fields.insert("when_to_use".into(), Value::from(when_to_use));
4030 fields.insert("instructions".into(), Value::from(instructions));
4031 if let Some(ns) = &ns {
4032 fields.insert("namespace".into(), Value::from(ns.clone()));
4033 }
4034 let existing = sub
4036 .grains_of_type(
4037 crate::model::grain_type::SKILL,
4038 ns.as_deref(),
4039 ReadOpts { live_only: true, since_ms: None },
4040 )
4041 .ok()?
4042 .into_iter()
4043 .find(|g| g.skill_name() == Some(name.as_str()))
4044 .map(|g| g.hash);
4045 Some(SkillFields { fields, name, n_steps: steps.len(), existing })
4046}
4047
4048fn derived_fact_fields(
4049 target: &TargetRef,
4050 relation: &str,
4051 object: &str,
4052 cited: &[String],
4053 ns_by_hash: &std::collections::BTreeMap<String, String>,
4054) -> Option<serde_json::Map<String, Value>> {
4055 if target.scheme() != "entity" {
4056 return None;
4057 }
4058 let subject = target
4059 .opaque()
4060 .rsplit_once('/')
4061 .map(|(_, s)| s)
4062 .unwrap_or(target.opaque());
4063 if subject.is_empty() {
4064 return None;
4065 }
4066 let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
4067 for h in cited {
4068 if let Some(ns) = ns_by_hash.get(h) {
4069 if !ns.is_empty() {
4070 *ns_counts.entry(ns.as_str()).or_default() += 1;
4071 }
4072 }
4073 }
4074 let lesson_ns = ns_counts
4075 .iter()
4076 .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
4077 .map(|(ns, _)| ns.to_string());
4078 let mut fields = serde_json::Map::new();
4079 fields.insert("subject".into(), Value::from(subject));
4080 fields.insert("relation".into(), Value::from(relation));
4081 fields.insert("object".into(), Value::from(object));
4082 if let Some(ns) = lesson_ns {
4086 fields.insert("namespace".into(), Value::from(ns));
4087 }
4088 Some(fields)
4089}
4090
4091fn latest_verdicts(p: &LoopPersisted) -> BTreeMap<String, String> {
4096 let mut out = BTreeMap::new();
4097 for (rec_hash, applied) in &p.applied {
4098 let latest = p
4099 .outcomes
4100 .get(rec_hash)
4101 .and_then(|v| v.iter().max_by_key(|o| o.measured_at_ms))
4102 .map(|o| o.verdict.clone());
4103 for h in &applied.created_hashes {
4104 out.insert(h.clone(), latest.clone().unwrap_or_else(|| "unmeasured".into()));
4105 }
4106 }
4107 out
4108}
4109
4110pub const NEAR_DUPLICATE_COSINE: f64 = 0.90;
4113pub const NEAR_DUPLICATE_JACCARD: f64 = 0.60;
4116const NEAR_DUPLICATE_CAP: usize = 8;
4118
4119pub(crate) fn near_duplicates_of<S: SubstrateRead + ?Sized>(
4124 sub: &S,
4125 subject: &str,
4126 namespace: Option<&str>,
4127 text: &str,
4128) -> Vec<crate::recommendation::NearDuplicate> {
4129 use crate::analyzers::duplicate_sweep::{jaccard, tokenize};
4130 if subject.is_empty() || text.trim().is_empty() {
4131 return Vec::new();
4132 }
4133 let Ok(facts) = sub.grains_of_type(
4134 crate::model::grain_type::FACT,
4135 namespace,
4136 ReadOpts { live_only: true, since_ms: None },
4137 ) else {
4138 return Vec::new();
4139 };
4140 let mine = sub.embed(text).ok().flatten();
4141 let my_tokens = tokenize(text);
4142 let mut out: Vec<crate::recommendation::NearDuplicate> = facts
4143 .iter()
4144 .filter(|f| f.fact_relation() == Some("lesson"))
4145 .filter(|f| f.fact_subject().is_some_and(|s| normalize_ident(s) == normalize_ident(subject)))
4146 .filter_map(|f| {
4147 let other = f.fact_object()?;
4148 let (score, method, floor) = match (&mine, sub.embed(other).ok().flatten()) {
4149 (Some(a), Some(b)) => (cosine(a, &b), "cosine", NEAR_DUPLICATE_COSINE),
4150 _ => (jaccard(&my_tokens, &tokenize(other)), "jaccard", NEAR_DUPLICATE_JACCARD),
4151 };
4152 (score >= floor).then(|| crate::recommendation::NearDuplicate {
4153 hash: f.hash.clone(),
4154 score: (score * 1000.0).round() / 1000.0,
4155 method: method.into(),
4156 })
4157 })
4158 .collect();
4159 out.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal).then(a.hash.cmp(&b.hash)));
4160 out.truncate(NEAR_DUPLICATE_CAP);
4161 out
4162}
4163
4164fn cosine(a: &[f32], b: &[f32]) -> f64 {
4165 if a.len() != b.len() || a.is_empty() {
4166 return 0.0;
4167 }
4168 let (mut dot, mut na, mut nb) = (0f64, 0f64, 0f64);
4169 for (x, y) in a.iter().zip(b) {
4170 dot += *x as f64 * *y as f64;
4171 na += *x as f64 * *x as f64;
4172 nb += *y as f64 * *y as f64;
4173 }
4174 if na == 0.0 || nb == 0.0 {
4175 0.0
4176 } else {
4177 dot / (na.sqrt() * nb.sqrt())
4178 }
4179}
4180
4181const CONSOLIDATION_INSTRUCTIONS: &str = " (8) {\"kind\":\"consolidation\",\"lesson\":\"...\",\
4185\"supersedes\":[\"<hash>\",...]} with the same entity target — ONLY in answer to a \
4186'Lesson pile' finding, which lists the live lessons on one entity that exceed \
4187its budget. Write ONE short imperative rule (max 240 chars) that says what \
4188those lessons say together, dropping nothing a lesson that measured 'held' \
4189required and keeping nothing only a lesson that measured 'regressed' or \
4190'drifted' added; 'supersedes' MUST be exactly the hashes that finding lists \
4191(cite them as evidence too). Applying it replaces every listed lesson with \
4192the one line; the reviewer can restore them all.";
4193
4194fn requires_gating(kind: ActionKind) -> bool {
4199 matches!(
4200 kind,
4201 ActionKind::CodeRevision | ActionKind::AdapterRevision
4202 )
4203}
4204
4205fn stamp(
4206 m: &AnalyzerManifest,
4207 params: &crate::manifest::Params,
4208 d: crate::recommendation::RecDraft,
4209 now_ms: i64,
4210) -> Result<Recommendation> {
4211 let target = TargetRef::parse(&d.target_ref)?;
4212 crate::recommendation::validate_code_rules(
4216 d.action_kind,
4217 target.target_class(),
4218 d.evalset_hash.as_deref(),
4219 )?;
4220 let revert_of = match (&d.action_kind, &d.proposal) {
4223 (ActionKind::Revert, Proposal::Data { data }) => {
4224 data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
4225 }
4226 _ => None,
4227 };
4228 let dedup = match revert_of.as_deref() {
4229 Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
4230 None => dedup_key(m.family(), &d.target_ref, d.action_kind),
4231 };
4232 let destructive = match &d.proposal {
4233 Proposal::Cal { cal } => cal::contains_destructive(cal),
4234 _ => false,
4235 };
4236 let rollbackable = match &d.proposal {
4237 Proposal::Cal { .. } => !destructive,
4238 Proposal::Edit { .. } => false,
4239 Proposal::Data { .. } => requires_gating(d.action_kind),
4243 };
4244 let mut evidence = d.evidence;
4245 evidence.truncate(MAX_EVIDENCE);
4246 let origin = match m.trust_class {
4251 crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
4252 _ => Origin::Builtin,
4253 };
4254 Ok(Recommendation {
4255 hash: String::new(),
4256 analyzer: m.id.clone(),
4257 params_snapshot: params.snapshot(),
4258 origin,
4259 target_ref: target.as_string(),
4260 action_kind: d.action_kind,
4261 dedup_key: dedup,
4262 summary: d.summary,
4263 severity: d.severity,
4264 proposal: d.proposal,
4265 destructive,
4266 rollbackable,
4267 evidence,
4268 evidence_query: d.evidence_query,
4269 metric: d.metric,
4270 confidence: d.confidence,
4271 importance: d.importance,
4272 created_at_ms: now_ms,
4273 guidance: None,
4274 evalset_hash: d.evalset_hash,
4275 near_duplicate_of: Vec::new(),
4276 replay: None,
4277 status: RecStatus::Pending,
4278 })
4279}
4280
4281fn validate_because(because: &str) -> Result<String> {
4282 let trimmed = because.trim();
4283 if trimmed.is_empty() {
4284 return Err(Error::InvalidProposal(
4285 "a BECAUSE reason is required".into(),
4286 ));
4287 }
4288 if trimmed.chars().count() > MAX_BECAUSE {
4289 return Err(Error::InvalidProposal(format!(
4290 "BECAUSE exceeds {MAX_BECAUSE} chars"
4291 )));
4292 }
4293 Ok(trimmed.to_string())
4294}
4295
4296fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
4297 for req in &m.requires {
4298 match req {
4299 Capability::Forks if !caps.forks => return Some("forks"),
4300 Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
4301 Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
4302 _ => {}
4303 }
4304 }
4305 None
4306}
4307
4308fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
4309 p.config.get(analyzer_id).and_then(|c| c.severity_floor)
4310}
4311
4312fn gate(
4313 opts: &RunOptions,
4314 p: &LoopPersisted,
4315 new_grains: u64,
4316 new_errors: u64,
4317 now_ms: i64,
4318) -> Option<SkipReason> {
4319 let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
4320 if !any {
4321 return None;
4322 }
4323 let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
4324 let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
4325 let stale_ok = opts
4326 .if_stale_ms
4327 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
4328 if min_new_ok || min_err_ok || stale_ok {
4329 return None;
4330 }
4331 if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
4333 Some(SkipReason::NotStale)
4334 } else {
4335 Some(SkipReason::MinNewNotMet)
4336 }
4337}
4338
4339#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
4341pub(crate) struct NewSince {
4342 pub grains: u64,
4344 pub error_events: u64,
4346 pub events: u64,
4348 pub sessions: u64,
4350}
4351
4352fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
4353 let opts = ReadOpts {
4354 live_only: false,
4355 since_ms: watermark.map(|w| w + 1),
4356 };
4357 let mut n = NewSince::default();
4358 let mut sessions: BTreeSet<&str> = BTreeSet::new();
4359 let mut events_held: Vec<GrainRecord> = Vec::new();
4360 for t in [
4361 crate::model::grain_type::FACT,
4362 crate::model::grain_type::EVENT,
4363 crate::model::grain_type::TOOL,
4364 crate::model::grain_type::OBSERVATION,
4365 ] {
4366 let g = sub.grains_of_type(t, None, opts)?;
4367 n.grains += g.len() as u64;
4368 if t == crate::model::grain_type::TOOL {
4370 n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
4371 }
4372 if t == crate::model::grain_type::EVENT {
4373 n.events = g.len() as u64;
4374 events_held = g;
4375 }
4376 }
4377 for e in &events_held {
4378 if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
4379 sessions.insert(sid);
4380 }
4381 }
4382 n.sessions = sessions.len() as u64;
4383 Ok(n)
4384}
4385
4386fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
4393 if !c.is_set() {
4394 return None;
4395 }
4396 let time_ok = c
4397 .every_ms
4398 .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
4399 let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
4400 let events_ok = c.every_events.is_some_and(|m| new.events >= m);
4401 let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
4402 if time_ok || grains_ok || events_ok || sessions_ok {
4403 None
4404 } else {
4405 Some(SkipReason::CadenceNotDue)
4406 }
4407}
4408
4409fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
4410 let grains = sub.grains_of_type(
4411 crate::model::grain_type::RECOMMENDATION,
4412 Some(LOOP_NS),
4413 ReadOpts {
4414 live_only: false,
4415 since_ms: None,
4416 },
4417 )?;
4418 let mut set = BTreeSet::new();
4419 for g in grains {
4420 let status = p
4421 .status_index
4422 .get(&g.hash)
4423 .copied()
4424 .unwrap_or(RecStatus::Pending);
4425 if matches!(
4432 status,
4433 RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
4434 ) {
4435 if let Some(key) = g.str_field("dedup_key") {
4436 set.insert(key.to_string());
4437 }
4438 }
4439 }
4440 Ok(set)
4441}
4442
4443pub(crate) fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
4450 const BASE_MS: i64 = 7 * 86_400_000;
4451 const CAP_MS: i64 = 90 * 86_400_000;
4452 let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
4453 let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
4454 *strikes = strikes.saturating_add(1);
4455 p.cooldowns.insert(dedup_key, now_ms + interval);
4456}
4457
4458fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
4459 let g = sub
4460 .grain(rec_hash)?
4461 .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
4462 Recommendation::from_fields(rec_hash, &g.fields)
4463}
4464
4465pub(crate) fn is_definition_statement(line: &str) -> bool {
4473 let up = line.trim_start().to_ascii_uppercase();
4474 up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
4475}
4476
4477const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
4480 Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
4481 approve it to acknowledge it and let it expire.";
4482
4483const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
4488 (evalset hash + run id + stats) — use apply_gated";
4489
4490const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
4491 guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
4492 acknowledge it and let it expire.";
4493
4494pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
4507 match proposal {
4508 Proposal::Cal { .. } => Ok(()),
4509 Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
4510 Proposal::Data { data } => {
4516 if requires_gating(action_kind)
4517 || data.get("revert_of").and_then(Value::as_str).is_some()
4518 {
4519 Ok(())
4520 } else {
4521 Err(Error::InvalidProposal(ADVISORY_DATA.into()))
4522 }
4523 }
4524 }
4525}
4526
4527#[cfg(test)]
4528mod definition_body_tests {
4529 use super::safe_definition_body;
4530
4531 #[test]
4532 fn ordinary_bodies_pass() {
4533 assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
4534 assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
4535 }
4536
4537 #[test]
4538 fn a_body_cannot_close_its_own_block_or_carry_destruction() {
4539 assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
4544 assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
4545 assert!(!safe_definition_body("RECALL facts FORGET abc"));
4547 assert!(!safe_definition_body("recall facts purge older than 1d"));
4548 assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
4549 assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
4551 assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
4553 }
4554}
4555
4556#[cfg(test)]
4557mod plan_edit_tests {
4558 use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
4559 use serde_json::json;
4560
4561 fn plan() -> serde_json::Value {
4562 json!({
4563 "nodes": ["fetch", "review", "post"],
4564 "edges": [
4565 {"src": "fetch", "dst": "review"},
4566 {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
4567 ],
4568 "bindings": {"fetch": "sha256:tool1"},
4569 "retries": {"fetch": 1}
4570 })
4571 }
4572
4573 #[test]
4574 fn the_allowlist_admits_thresholds_and_refuses_topology() {
4575 assert!(plan_edit_allowed("edges.1.cond"));
4576 assert!(plan_edit_allowed("edges.1.max_cycles"));
4577 assert!(plan_edit_allowed("retries.fetch"));
4578 for path in [
4581 "nodes",
4582 "nodes.0",
4583 "edges.0.src",
4584 "edges.0.dst",
4585 "edges",
4586 "bindings.fetch",
4587 "edges.x.cond",
4588 "",
4589 ] {
4590 assert!(!plan_edit_allowed(path), "{path} must not be editable");
4591 }
4592 }
4593
4594 #[test]
4595 fn values_are_type_checked_against_the_field() {
4596 assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
4599 assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
4600 assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
4601 assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
4602 assert!(plan_value_ok("retries.fetch", &json!(3)));
4603 assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
4604 assert!(!plan_value_ok("edges.1.cond", &json!(" ")));
4605 assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
4606 assert!(!plan_value_ok("edges.0.src", &json!("other")));
4607 }
4608
4609 #[test]
4610 fn get_reads_through_arrays_and_objects_and_absence_is_null() {
4611 let p = plan();
4612 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
4613 assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
4614 assert_eq!(plan_get(&p, "retries.review"), json!(null));
4617 assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
4618 assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
4619 }
4620
4621 #[test]
4622 fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
4623 let mut p = plan();
4624 assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
4625 assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
4626 assert!(plan_set(&mut p, "retries.review", json!(2)));
4627 assert_eq!(plan_get(&p, "retries.review"), json!(2));
4628 assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
4629 assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
4630 let mut q = plan();
4635 q.as_object_mut().unwrap().remove("retries");
4636 assert!(plan_set(&mut q, "retries.greet", json!(1)));
4637 assert_eq!(plan_get(&q, "retries.greet"), json!(1));
4638 let mut r = plan();
4640 r.as_object_mut().unwrap().remove("edges");
4641 assert!(!plan_set(&mut r, "edges.0.cond", json!("x")));
4642 }
4643}
4644
4645#[cfg(test)]
4646mod definition_proposal_tests {
4647 use super::is_definition_statement;
4648
4649 #[test]
4650 fn definition_statements_are_recognized_in_both_spellings() {
4651 assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
4652 assert!(is_definition_statement(" define template foo AS { x }"));
4653 assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
4654 assert!(!is_definition_statement("ADD fact {}"));
4656 assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
4657 assert!(!is_definition_statement("FORGET abc"));
4658 assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
4662 }
4663}