Skip to main content

areev_loop/
engine.rs

1//! The engine: the analyze → DISCOVER → ENRICH → validate/dedup → store
2//! pipeline, the run-outcome contract, and the review/apply/rollback lifecycle
3//! with the governance gates. The DETERMINISTIC output is a pure function of
4//! (store state, params, now); the optional LLM stages (§9) only *add* cited
5//! drafts (origin=llm, never auto-apply) and whitelisted guidance — with no
6//! backend they are the identity, so the deterministic path is unchanged.
7//! Auto-apply execution is gated behind a conservative shape check and stays
8//! off by default.
9
10use 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,
19};
20use crate::substrate::{Capabilities, OmsSubstrate, ReadOpts, SubstrateRead};
21use serde::{Deserialize, Serialize};
22use serde_json::{Map, Value};
23use std::collections::{BTreeMap, BTreeSet};
24
25/// The namespace the loop's own grains (recommendations, audit) live in.
26pub const LOOP_NS: &str = "areev-loop";
27
28/// Host-granted authority, per connection. `admin` implies all.
29#[derive(Debug, Clone, Copy, PartialEq, Eq)]
30pub enum Scope {
31    Read,
32    Write,
33    Review,
34    Apply,
35    Admin,
36}
37
38/// A set of granted scopes.
39#[derive(Debug, Clone, Default)]
40pub struct ScopeSet(Vec<Scope>);
41
42impl ScopeSet {
43    pub fn of(scopes: &[Scope]) -> Self {
44        ScopeSet(scopes.to_vec())
45    }
46    /// The local root of trust: whoever can run against the file holds all
47    /// scopes (the CLI/embedded posture).
48    pub fn all() -> Self {
49        ScopeSet(vec![Scope::Admin])
50    }
51    pub fn has(&self, s: Scope) -> bool {
52        self.0.contains(&Scope::Admin) || self.0.contains(&s)
53    }
54}
55
56/// A review decision.
57#[derive(Debug, Clone, Copy, PartialEq, Eq)]
58pub enum Decision {
59    Approve,
60    Reject,
61}
62
63/// Gating and scoping options for a run.
64#[derive(Debug, Clone, Default)]
65pub struct RunOptions {
66    pub min_new: Option<u64>,
67    pub min_new_errors: Option<u64>,
68    pub if_stale_ms: Option<i64>,
69    /// Optional global namespace filter (empty = all).
70    pub namespaces: Vec<String>,
71    /// Re-analyze the WHOLE memory this pass, not just grains since the last-run
72    /// watermark (`areev loop reflect`). Dedup/cooldowns still suppress anything
73    /// already queued, and the watermark still advances at the end — so a sweep
74    /// is safe to run any time and later runs stay incremental. Mainly widens the
75    /// watermark-sensitive inputs (tool-failure window, the non-parasitic LLM
76    /// evidence bundle) to the full history.
77    pub full_sweep: bool,
78    /// The principal that invoked this run. Recorded as co-creator on every
79    /// non-`Builtin` (LLM / external-command) recommendation the run stores,
80    /// so the review gate's self-approval block also fires for whoever
81    /// triggered the model that authored the finding — an LLM draft is
82    /// authored *via* its trigger, unlike a deterministic finding, which is
83    /// computed. `None` (a headless/scheduled run) records no co-creator.
84    pub triggering_actor: Option<String>,
85}
86
87/// Whether a run executed or was skipped.
88#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(rename_all = "lowercase")]
90pub enum RunOutcome {
91    Ran,
92    Skipped,
93}
94
95/// Why a run was a no-op. `LockHeld` is produced by the host adapter (a
96/// concurrent writer), surfaced here for a single contract.
97#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
98#[serde(rename_all = "snake_case")]
99pub enum SkipReason {
100    MinNewNotMet,
101    NotStale,
102    LockHeld,
103}
104
105/// One analyzer that did not contribute drafts, with why.
106#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
107pub struct AnalyzerSkip {
108    pub id: String,
109    pub reason: String,
110}
111
112/// The run-outcome contract (proposal §13): one shape across CLI/API/MCP/bindings.
113#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
114pub struct RunResult {
115    pub outcome: RunOutcome,
116    #[serde(skip_serializing_if = "Option::is_none")]
117    pub skip_reason: Option<SkipReason>,
118    pub new_grains: u64,
119    pub new_error_events: u64,
120    pub proposed: u64,
121    pub deduped: u64,
122    pub stored: u64,
123    /// Of the stored recommendations, how many were auto-applied by policy.
124    #[serde(default)]
125    pub auto_applied: u64,
126    #[serde(default)]
127    pub analyzers_run: Vec<String>,
128    #[serde(default)]
129    pub analyzers_skipped: Vec<AnalyzerSkip>,
130}
131
132impl RunResult {
133    fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
134        RunResult {
135            outcome: RunOutcome::Skipped,
136            skip_reason: Some(reason),
137            new_grains,
138            new_error_events,
139            proposed: 0,
140            deduped: 0,
141            stored: 0,
142            auto_applied: 0,
143            analyzers_run: vec![],
144            analyzers_skipped: vec![],
145        }
146    }
147
148    pub fn ran(&self) -> bool {
149        self.outcome == RunOutcome::Ran
150    }
151}
152
153/// The engine holds the registered analyzers, the host policy, and an optional
154/// LLM enrichment backend (§9).
155pub struct Engine {
156    analyzers: Vec<Box<dyn Analyzer>>,
157    policy: crate::policy::Policy,
158    /// Optional LLM backend. `None` → the DISCOVER/ENRICH stages are the
159    /// identity, so the pipeline is byte-for-byte the deterministic path.
160    llm: Option<Box<dyn crate::llm::LlmBackend>>,
161    /// Optional separate backend for the GROUND stage (§5.2, §11). `None` →
162    /// grounding rides `llm`. Lets a team point entailment at a cheaper or
163    /// specialized model (or take the generative model out of grounding
164    /// entirely) without changing the proposer/verifier.
165    ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
166}
167
168struct AnalysisPass {
169    survivors: Vec<Recommendation>,
170    proposed: u64,
171    deduped: u64,
172    analyzers_run: Vec<String>,
173    analyzers_skipped: Vec<AnalyzerSkip>,
174}
175
176impl Engine {
177    /// An engine with the default built-ins and a default (fully closed)
178    /// policy — nothing auto-applies, no LLM.
179    pub fn with_builtins() -> Self {
180        Engine {
181            analyzers: crate::analyzer::builtin_analyzers(),
182            policy: crate::policy::Policy::default(),
183            llm: None,
184            ground_llm: None,
185        }
186    }
187
188    /// An engine with no analyzers (register your own).
189    pub fn empty() -> Self {
190        Engine {
191            analyzers: vec![],
192            policy: crate::policy::Policy::default(),
193            llm: None,
194            ground_llm: None,
195        }
196    }
197
198    /// Install a host policy (the only place auto-apply is granted).
199    pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
200        self.policy = policy;
201        self
202    }
203
204    /// Attach an optional LLM enrichment backend (§9). Only ever *adds* cited
205    /// draft recommendations (stamped `origin = llm`, never auto-applied) and
206    /// whitelisted guidance notes — it can never gate or rewrite deterministic
207    /// output.
208    pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
209        self.llm = Some(backend);
210        self
211    }
212
213    /// Attach a separate backend for the GROUND stage (§5.2). Without this,
214    /// grounding uses the `with_llm` backend. Independent of the proposer so an
215    /// operator can run entailment on a cheaper/specialized model.
216    pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
217        self.ground_llm = Some(backend);
218        self
219    }
220
221    pub fn policy(&self) -> &crate::policy::Policy {
222        &self.policy
223    }
224
225    /// Register an additional analyzer (the linked-Rust seam).
226    pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
227        self.analyzers.push(analyzer);
228    }
229
230    pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
231        &self.analyzers
232    }
233
234    /// Run the exact production analysis/validation path without Phase 0
235    /// measurement or Phase 3 persistence. The immutable substrate borrow is
236    /// the replay safety boundary: recommendations, audit grains, state,
237    /// cooldowns, outcomes, and the op-log cannot be changed here.
238    ///
239    /// `overrides` is keyed by full analyzer id and overlays the file's stored
240    /// parameter map. Unknown keys fail closed through `resolve_params` and
241    /// surface as an analyzer skip, exactly as in a production run.
242    pub fn analyze_only<S: OmsSubstrate>(
243        &self,
244        sub: &S,
245        opts: &RunOptions,
246        overrides: &BTreeMap<String, Map<String, Value>>,
247        now_ms: i64,
248    ) -> Result<Vec<Recommendation>> {
249        let persisted = LoopPersisted::from_value(sub.load_state()?)?;
250        let analysis_watermark = if opts.full_sweep {
251            None
252        } else {
253            persisted.state.watermark_ms
254        };
255        Ok(self
256            .analysis_pass(
257                sub,
258                &persisted,
259                opts,
260                overrides,
261                analysis_watermark,
262                now_ms,
263                &[],
264            )?
265            .survivors)
266    }
267
268    /// Run one analysis pass. Idempotent under `dedup_key`; the watermark is
269    /// advanced at the end, so a crashed run simply re-runs.
270    pub fn run<S: OmsSubstrate>(
271        &self,
272        sub: &mut S,
273        opts: &RunOptions,
274        now_ms: i64,
275    ) -> Result<RunResult> {
276        let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
277        let watermark = persisted.state.watermark_ms;
278        // A full sweep analyzes the whole memory (watermark ignored for the
279        // analysis inputs), while gating, `new` counts, and the end-of-run
280        // watermark advance still use the real watermark. Dedup/cooldowns keep
281        // it from re-proposing what is already queued.
282        let analysis_watermark = if opts.full_sweep { None } else { watermark };
283
284        let (new_grains, new_error_events) = count_new(sub, watermark)?;
285        if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
286            return Ok(RunResult::skipped(reason, new_grains, new_error_events));
287        }
288
289        // Phase 0: re-measure applied recommendations due for review (the
290        // Verify gate). Records a measured outcome per due recommendation.
291        let outcome_inputs = measure_outcomes(sub, &mut persisted, now_ms)?;
292
293        let AnalysisPass {
294            survivors,
295            proposed,
296            deduped,
297            analyzers_run,
298            analyzers_skipped,
299        } = self.analysis_pass(
300            &*sub,
301            &persisted,
302            opts,
303            &BTreeMap::new(),
304            analysis_watermark,
305            now_ms,
306            &outcome_inputs,
307        )?;
308
309        // Phase 3 (needs &mut): store survivors + propose audit, then
310        // auto-apply the ones the host policy grants (all gates in §6.3).
311        let mut stored = 0u64;
312        let mut auto_applied = 0u64;
313        for mut rec in survivors {
314            let spec = rec.to_grain_spec(LOOP_NS)?;
315            let hash = sub.put_grain(&spec)?;
316            rec.hash = hash.clone();
317            let actor = format!("engine:{}", rec.analyzer);
318            let audit = AuditRecord {
319                rec_hash: hash.clone(),
320                from: None,
321                to: RecStatus::Pending,
322                actor: actor.clone(),
323                observer_type: ObserverType::System,
324                because: "analyzer proposed".into(),
325                previous_audit_hash: None,
326                gating: None,
327                at_ms: now_ms,
328            };
329            let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
330            persisted
331                .status_index
332                .insert(hash.clone(), RecStatus::Pending);
333            persisted.creators.insert(hash.clone(), actor);
334            // An LLM or external-command finding exists because someone ran
335            // it — record that principal too, so review can refuse the
336            // trigger approving their own model's output. Builtin analyzers
337            // stay engine-only: deterministic output has no human author.
338            if !matches!(rec.origin, Origin::Builtin) {
339                if let Some(trigger) = &opts.triggering_actor {
340                    persisted.co_creators.insert(hash.clone(), trigger.clone());
341                }
342            }
343            persisted.audit_heads.insert(hash.clone(), audit_hash);
344            stored += 1;
345
346            if self.can_auto_apply(&*sub, &rec) {
347                self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
348                auto_applied += 1;
349            }
350        }
351
352        persisted.state.last_run_ms = Some(now_ms);
353        persisted.state.watermark_ms = Some(now_ms);
354        sub.store_state(&persisted.to_value()?)?;
355
356        Ok(RunResult {
357            outcome: RunOutcome::Ran,
358            skip_reason: None,
359            new_grains,
360            new_error_events,
361            proposed,
362            deduped,
363            stored,
364            auto_applied,
365            analyzers_run,
366            analyzers_skipped,
367        })
368    }
369
370    /// Shared Phase 1–2 implementation for production and replay. Cooldowns
371    /// and the live recommendation queue are honored in both modes; only the
372    /// production caller proceeds into the mutating Phase 3 below.
373    #[allow(clippy::too_many_arguments)]
374    fn analysis_pass<S: OmsSubstrate>(
375        &self,
376        sub: &S,
377        persisted: &LoopPersisted,
378        opts: &RunOptions,
379        external_overrides: &BTreeMap<String, Map<String, Value>>,
380        analysis_watermark: Option<i64>,
381        now_ms: i64,
382        outcome_inputs: &[OutcomeInput],
383    ) -> Result<AnalysisPass> {
384        let existing = existing_dedup_keys(sub, persisted)?;
385        let mut analyzers_run = Vec::new();
386        let mut analyzers_skipped = Vec::new();
387        let mut candidates: Vec<Recommendation> = Vec::new();
388        let caps = sub.capabilities();
389
390        for analyzer in &self.analyzers {
391            let m = analyzer.manifest();
392            let cfg = persisted.config.get(&m.id);
393            let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
394            if !enabled {
395                analyzers_skipped.push(AnalyzerSkip {
396                    id: m.id.clone(),
397                    reason: "disabled".into(),
398                });
399                continue;
400            }
401            if self.policy.denies(m.family()) {
402                analyzers_skipped.push(AnalyzerSkip {
403                    id: m.id.clone(),
404                    reason: "denied by host policy".into(),
405                });
406                continue;
407            }
408            if let Some(missing) = missing_capability(m, caps) {
409                analyzers_skipped.push(AnalyzerSkip {
410                    id: m.id.clone(),
411                    reason: format!("missing capability: {missing}"),
412                });
413                continue;
414            }
415            let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
416            if let Some(extra) = external_overrides.get(&m.id) {
417                for (key, value) in extra {
418                    param_overrides.insert(key.clone(), value.clone());
419                }
420            }
421            let params = match m.resolve_params(&param_overrides) {
422                Ok(p) => p,
423                Err(e) => {
424                    analyzers_skipped.push(AnalyzerSkip {
425                        id: m.id.clone(),
426                        reason: e.to_string(),
427                    });
428                    continue;
429                }
430            };
431            let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
432            let ns_slice: &[String] = if ns_owned.is_empty() {
433                &opts.namespaces
434            } else {
435                &ns_owned
436            };
437            let reader: &dyn SubstrateRead = sub;
438            let ctx = AnalyzeCtx::new(
439                reader,
440                &params,
441                ns_slice,
442                analysis_watermark,
443                now_ms,
444                outcome_inputs,
445            );
446            match analyzer.analyze(&ctx) {
447                Ok(drafts) => {
448                    analyzers_run.push(m.id.clone());
449                    for draft in drafts {
450                        match stamp(m, &params, draft, now_ms) {
451                            Ok(rec) => candidates.push(rec),
452                            Err(e) => analyzers_skipped.push(AnalyzerSkip {
453                                id: m.id.clone(),
454                                reason: e.to_string(),
455                            }),
456                        }
457                    }
458                }
459                Err(e) => analyzers_skipped.push(AnalyzerSkip {
460                    id: m.id.clone(),
461                    reason: e.to_string(),
462                }),
463            }
464        }
465
466        if self.llm.is_some() {
467            candidates.extend(self.discover(
468                sub,
469                &candidates,
470                analysis_watermark,
471                &opts.namespaces,
472                now_ms,
473            ));
474        }
475
476        let proposed = candidates.len() as u64;
477        let mut seen = BTreeSet::new();
478        let mut survivors = Vec::new();
479        for candidate in candidates {
480            let family = crate::manifest::analyzer_family(&candidate.analyzer);
481            let floor = [
482                severity_floor_for(persisted, &candidate.analyzer),
483                self.policy.severity_floor(family),
484            ]
485            .into_iter()
486            .flatten()
487            .max();
488            if floor.is_some_and(|floor| candidate.severity < floor) {
489                continue;
490            }
491            if !seen.insert(candidate.dedup_key.clone()) {
492                continue;
493            }
494            if existing.contains(&candidate.dedup_key) {
495                continue;
496            }
497            if persisted
498                .cooldowns
499                .get(&candidate.dedup_key)
500                .is_some_and(|until| now_ms < *until)
501            {
502                continue;
503            }
504            survivors.push(candidate);
505        }
506        let deduped = proposed - survivors.len() as u64;
507        if self.llm.is_some() {
508            self.enrich(&mut survivors);
509        }
510        Ok(AnalysisPass {
511            survivors,
512            proposed,
513            deduped,
514            analyzers_run,
515            analyzers_skipped,
516        })
517    }
518
519    /// DISCOVER (§9): ask the LLM for additional draft recommendations, given
520    /// the deterministic findings as *context* and a bounded, provenance-tagged
521    /// evidence bundle. Every returned draft must cite evidence present in the
522    /// bundle and target a memory/query surface; it is stamped `origin = llm`
523    /// (so it can never auto-apply) and enters the ordinary dedup/store path. A
524    /// failed or garbled response yields no drafts — never a failed run.
525    fn discover<S: OmsSubstrate>(
526        &self,
527        sub: &S,
528        candidates: &[Recommendation],
529        watermark: Option<i64>,
530        namespaces: &[String],
531        now_ms: i64,
532    ) -> Vec<Recommendation> {
533        let Some(llm) = &self.llm else {
534            return Vec::new();
535        };
536        let findings: Vec<crate::llm::FindingBrief> = candidates
537            .iter()
538            .take(32)
539            .map(|c| crate::llm::FindingBrief {
540                analyzer: c.analyzer.clone(),
541                summary: c.summary.render(),
542                target: c.target_ref.clone(),
543                severity: c.severity.as_str().to_string(),
544            })
545            .collect();
546        // Evidence bundle. Seeded first from the grains the deterministic
547        // findings cite, THEN — the non-parasitic step (§11) — topped up with
548        // RECENT grains (created since the last run) so the LLM gets its own
549        // lens and can find issues in grains no analyzer flagged. Without this
550        // the LLM could only elaborate near what determinism already caught.
551        let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
552        let mut bundle: BTreeSet<String> = BTreeSet::new();
553        for c in candidates {
554            for h in &c.evidence {
555                if !bundle.contains(h) {
556                    if let Ok(Some(g)) = sub.grain(h) {
557                        push_evidence(&mut evidence, &mut bundle, &g);
558                    }
559                }
560            }
561        }
562        let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
563            vec![None]
564        } else {
565            namespaces.iter().map(|n| Some(n.as_str())).collect()
566        };
567        let opts = ReadOpts { live_only: true, since_ms: watermark };
568        'seed: for gt in [
569            crate::model::grain_type::FACT,
570            crate::model::grain_type::OBSERVATION,
571        ] {
572            for ns in &scan_ns {
573                if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
574                    for g in recent {
575                        if evidence.len() >= 64 {
576                            break 'seed;
577                        }
578                        push_evidence(&mut evidence, &mut bundle, &g);
579                    }
580                }
581            }
582        }
583        if evidence.is_empty() {
584            return Vec::new(); // nothing to reflect on
585        }
586        // PROPOSE (§5.1): the abstention-legitimate objective — "nothing to
587        // report" is a first-class, zero-penalty answer. The operator-taste
588        // history (recent approve/reject decisions on llm findings) is passed so
589        // the model learns what this reviewer accepts.
590        let (approved, rejected) = self.llm_history(sub);
591        let request = crate::llm::LlmRequest {
592            loop_proto: 1,
593            op: "discover",
594            instructions: DISCOVER_INSTRUCTIONS,
595            findings: findings.clone(),
596            evidence: evidence.clone(),
597            rejected,
598            approved,
599        };
600        let Ok(body) = serde_json::to_string(&request) else {
601            return Vec::new();
602        };
603        let raw = match llm.complete(&body) {
604            Ok(r) => r,
605            Err(_) => return Vec::new(), // fail-soft
606        };
607        // Cheap structural validation (cite-check + target class); collect the
608        // survivors for the verifier. Storing the normalized target string
609        // avoids a TargetRef clone through the pipeline.
610        let mut validated: Vec<(crate::llm::LlmDraft, String, Vec<String>)> = Vec::new();
611        for d in crate::llm::parse_discover(&raw)
612            .recommendations
613            .into_iter()
614            .take(crate::llm::MAX_LLM_DRAFTS)
615        {
616            let cited: Vec<String> =
617                d.evidence.iter().filter(|h| bundle.contains(*h)).cloned().collect();
618            if cited.is_empty() {
619                continue; // uncited → drop (no fabrication)
620            }
621            let Ok(target) = TargetRef::parse(&d.target) else {
622                continue;
623            };
624            let tc = target.target_class();
625            if tc != "memory" && tc != "query" {
626                continue; // never prompt/host
627            }
628            validated.push((d, target.as_string(), cited));
629        }
630        if validated.is_empty() {
631            return Vec::new();
632        }
633        // GROUND → VERIFY → ROUTE (§5.2–5.4): only drafts that survive an
634        // independent grounding entailment check *and* an adversarial
635        // verification pass (each a separate call — proposer ≠ scorer) reach the
636        // queue, stamped with the verifier's calibrated confidence.
637        // GROUND may run on a separate backend (§11); VERIFY always uses the
638        // main llm (the proposer≠scorer independence is on VERIFY, not GROUND).
639        let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
640        self.verify_drafts(&**llm, ground, &validated, &evidence, now_ms)
641    }
642
643    /// GROUND → VERIFY → ROUTE (§5.2–5.4). Two independent model calls, batched
644    /// over the drafts: a grounding-entailment gate ("does the cited evidence
645    /// support the claim?"), then an adversarial keep/kill with a calibrated
646    /// confidence. A draft reaches the queue only if it is grounded **and** kept
647    /// **and** clears the confidence floor. Any failed call drops the whole LLM
648    /// contribution for the run (safe default), never the run.
649    fn verify_drafts(
650        &self,
651        llm: &dyn crate::llm::LlmBackend,
652        ground: &dyn crate::llm::LlmBackend,
653        validated: &[(crate::llm::LlmDraft, String, Vec<String>)],
654        evidence: &[crate::llm::EvidenceItem],
655        now_ms: i64,
656    ) -> Vec<Recommendation> {
657        use crate::llm::*;
658        let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
659            evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
660        let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
661            cited
662                .iter()
663                .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
664                .collect()
665        };
666
667        // GROUND (§5.2): decompose-then-entail per draft, batched into one call.
668        let claims: Vec<GroundItem> = validated
669            .iter()
670            .enumerate()
671            .map(|(i, (d, _t, cited))| GroundItem {
672                id: i,
673                claim: cap(&d.summary, MAX_SUMMARY_LEN),
674                evidence: ev_for(cited),
675            })
676            .collect();
677        let ground_req = GroundRequest {
678            loop_proto: 1,
679            op: "ground",
680            instructions: GROUND_INSTRUCTIONS,
681            claims,
682        };
683        let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
684            .ok()
685            .and_then(|b| ground.complete(&b).ok())
686        {
687            Some(raw) => parse_ground(&raw)
688                .results
689                .into_iter()
690                .filter(|r| r.supported)
691                .map(|r| r.id)
692                .collect(),
693            None => return Vec::new(),
694        };
695        if grounded.is_empty() {
696            return Vec::new();
697        }
698
699        // VERIFY (§5.3): adversarial keep/kill over the grounded drafts, a
700        // separate call from the proposer. Soundness + abstention only — NOT
701        // novelty. Novelty is steered at DISCOVER and settled by human review;
702        // asking a weak verifier to judge it just makes it hallucinate "already
703        // known" and kill genuine findings (§11).
704        let items: Vec<VerifyItem> = validated
705            .iter()
706            .enumerate()
707            .filter(|(i, _)| grounded.contains(i))
708            .map(|(i, (d, t, cited))| VerifyItem {
709                id: i,
710                summary: cap(&d.summary, MAX_SUMMARY_LEN),
711                target: t.clone(),
712                evidence: ev_for(cited),
713            })
714            .collect();
715        let verify_req = VerifyRequest {
716            loop_proto: 1,
717            op: "verify",
718            instructions: VERIFY_INSTRUCTIONS,
719            findings: items,
720        };
721        let verdicts: std::collections::BTreeMap<usize, f64> =
722            match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
723                Some(raw) => parse_verify(&raw)
724                    .results
725                    .into_iter()
726                    .filter(|r| r.keep)
727                    .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
728                    .collect(),
729                None => return Vec::new(),
730            };
731
732        // ROUTE (§5.4): grounded ∧ kept ∧ verifier-confidence ≥ floor. The
733        // verifier's confidence (the independent signal) is what we trust and
734        // stamp — not the proposer's self-report.
735        let mut out = Vec::new();
736        for (i, (d, target_str, cited)) in validated.iter().enumerate() {
737            if let Some(&conf) = verdicts.get(&i) {
738                if conf >= MIN_LLM_CONFIDENCE {
739                    out.push(stamp_llm(
740                        llm.model(),
741                        d,
742                        target_str.clone(),
743                        cited.clone(),
744                        conf,
745                        now_ms,
746                    ));
747                }
748            }
749        }
750        out
751    }
752
753    /// Recent operator decisions on `origin = llm` findings — approved (incl.
754    /// applied) and rejected summaries, most-recent first and bounded — so
755    /// DISCOVER can learn what this reviewer accepts (§9). Best-effort: a read
756    /// failure yields empty history, never an error.
757    fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
758        const MAX: usize = 20;
759        let Ok(mut recs) = self.recommendations(sub, None) else {
760            return (Vec::new(), Vec::new());
761        };
762        recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
763        recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
764        let mut approved = Vec::new();
765        let mut rejected = Vec::new();
766        for r in &recs {
767            match r.status {
768                RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
769                    if approved.len() < MAX =>
770                {
771                    approved.push(r.summary.render());
772                }
773                RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
774                _ => {}
775            }
776        }
777        (approved, rejected)
778    }
779
780    /// ENRICH (§9): ask the LLM to add a short guidance note to the surviving
781    /// deterministic recommendations. Whitelist-only — only `guidance` is
782    /// merged (capped), and only onto recs that don't already have one; the
783    /// engine-templated summary is never touched. Fail-soft.
784    fn enrich(&self, survivors: &mut [Recommendation]) {
785        let Some(llm) = &self.llm else {
786            return;
787        };
788        if survivors.is_empty() {
789            return;
790        }
791        let findings: Vec<crate::llm::FindingBrief> = survivors
792            .iter()
793            .map(|r| crate::llm::FindingBrief {
794                analyzer: r.analyzer.clone(),
795                summary: r.summary.render(),
796                target: r.target_ref.clone(),
797                severity: r.severity.as_str().to_string(),
798            })
799            .collect();
800        let request = crate::llm::LlmRequest {
801            loop_proto: 1,
802            op: "enrich",
803            instructions: ENRICH_INSTRUCTIONS,
804            findings,
805            evidence: Vec::new(),
806            rejected: Vec::new(),
807            approved: Vec::new(),
808        };
809        let Ok(body) = serde_json::to_string(&request) else {
810            return;
811        };
812        let raw = match llm.complete(&body) {
813            Ok(r) => r,
814            Err(_) => return,
815        };
816        for note in crate::llm::parse_enrich(&raw).notes {
817            if note.guidance.trim().is_empty() {
818                continue;
819            }
820            if let Some(r) = survivors
821                .iter_mut()
822                .find(|r| r.target_ref == note.target && r.guidance.is_none())
823            {
824                r.guidance = Some(crate::llm::cap(&note.guidance, crate::llm::MAX_GUIDANCE_LEN));
825            }
826        }
827    }
828
829    /// Evaluate the auto-apply gate (§6.3) — ALL preconditions must hold:
830    /// host opt-in + policy grant, builtin origin, memory/query target,
831    /// non-destructive, and engine-side shape verification: SUPERSEDE-only
832    /// structural curation (never an ADD that introduces evidence-derived
833    /// text) whose every replacement is **value-identical** to the grain it
834    /// supersedes (the exact-equality check — a near-duplicate consolidation
835    /// stays pending). A default (closed) policy never grants, so nothing
836    /// auto-applies.
837    fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
838        if !rec.origin.auto_apply_eligible() || rec.destructive {
839            return false;
840        }
841        // The analyzer must declare its curation auto-appliable. An analyzer
842        // whose manifest is `Never` (e.g. fork surfacing — a lossy merge) is
843        // never auto-applied even if the payload passes the shape check.
844        let manifest_ok = self
845            .analyzers
846            .iter()
847            .map(|a| a.manifest())
848            .find(|m| m.id == rec.analyzer)
849            .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
850        if !manifest_ok {
851            return false;
852        }
853        let Ok(target) = TargetRef::parse(&rec.target_ref) else {
854            return false;
855        };
856        let family = crate::manifest::analyzer_family(&rec.analyzer);
857        if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
858            return false;
859        }
860        // Shape verification: only a CAL batch of pure SUPERSEDE statements
861        // whose replacements change no value is structural curation. An ADD
862        // (introducing content), a FORGET (destructive), or a supersession
863        // that alters any field disqualifies.
864        match &rec.proposal {
865            Proposal::Cal { cal } => cal
866                .lines()
867                .map(str::trim)
868                .filter(|l| !l.is_empty())
869                .all(|l| supersede_is_value_identical(sub, l)),
870            _ => false,
871        }
872    }
873
874    /// Apply a recommendation as `policy:auto` (the only `pending → applied`
875    /// path). Records the applied inverse + a hash-chained audit grain.
876    fn auto_apply<S: OmsSubstrate>(
877        &self,
878        sub: &mut S,
879        p: &mut LoopPersisted,
880        rec: &Recommendation,
881        now_ms: i64,
882    ) -> Result<()> {
883        let mut created = Vec::new();
884        if let Proposal::Cal { cal } = &rec.proposal {
885            // Belt and braces over the policy: `grants_auto_apply` already
886            // excludes the `query` class, so a definition rewrite cannot reach
887            // this path. If one ever did, it would apply with no recorded
888            // inverse and no human BECAUSE — refuse instead.
889            if cal.lines().map(str::trim).any(is_definition_statement) {
890                return Err(Error::InvalidProposal(
891                    "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
892                     auto-applied: it changes what every future context contains, so it \
893                     requires a human APPROVE + APPLY with BECAUSE"
894                        .into(),
895                ));
896            }
897            for r in sub.execute_cal(cal)? {
898                if let Some(h) = r.get("hash").and_then(Value::as_str) {
899                    created.push(h.to_string());
900                }
901            }
902        }
903        let applied = AppliedRecord {
904            applied_at_ms: now_ms,
905            target_ref: rec.target_ref.clone(),
906            rollbackable: rec.rollbackable,
907            created_hashes: created,
908            inverse_cal: None,
909            metric: rec.metric.clone(),
910        };
911        let prev = p.audit_heads.get(&rec.hash).cloned();
912        let audit = AuditRecord {
913            rec_hash: rec.hash.clone(),
914            from: Some(RecStatus::Pending),
915            to: RecStatus::Applied,
916            actor: "policy:auto".into(),
917            observer_type: ObserverType::Policy,
918            because: "auto-applied per host policy".into(),
919            previous_audit_hash: prev,
920            gating: None,
921            at_ms: now_ms,
922        };
923        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
924        p.audit_heads.insert(rec.hash.clone(), audit_hash);
925        p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
926        p.applied.insert(rec.hash.clone(), applied);
927        Ok(())
928    }
929
930    /// Approve or reject a pending recommendation. Requires the `review` scope,
931    /// a mandatory BECAUSE, and blocks self-approval against the creating actor.
932    #[allow(clippy::too_many_arguments)]
933    pub fn review<S: OmsSubstrate>(
934        &self,
935        sub: &mut S,
936        rec_hash: &str,
937        decision: Decision,
938        actor: &str,
939        observer: ObserverType,
940        scopes: &ScopeSet,
941        because: &str,
942        now_ms: i64,
943    ) -> Result<()> {
944        if !scopes.has(Scope::Review) {
945            return Err(Error::ScopeDenied("review".into()));
946        }
947        let because = validate_because(because)?;
948        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
949        let status = *p
950            .status_index
951            .get(rec_hash)
952            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
953        let to = match decision {
954            Decision::Approve => RecStatus::Approved,
955            Decision::Reject => RecStatus::Rejected,
956        };
957        if !status.can_transition_to(to, false) {
958            return Err(Error::LifecycleViolation(format!(
959                "{} -> {}",
960                status.as_str(),
961                to.as_str()
962            )));
963        }
964        if to == RecStatus::Approved {
965            if let Some(creator) = p.creators.get(rec_hash) {
966                if creator == actor {
967                    return Err(Error::SelfApproval(format!(
968                        "{actor} created this recommendation"
969                    )));
970                }
971            }
972            if let Some(trigger) = p.co_creators.get(rec_hash) {
973                if trigger == actor {
974                    return Err(Error::SelfApproval(format!(
975                        "{actor} triggered the run that authored this recommendation"
976                    )));
977                }
978            }
979        }
980        let prev = p.audit_heads.get(rec_hash).cloned();
981        let audit = AuditRecord {
982            rec_hash: rec_hash.into(),
983            from: Some(status),
984            to,
985            actor: actor.into(),
986            observer_type: observer,
987            because,
988            previous_audit_hash: prev,
989            gating: None,
990            at_ms: now_ms,
991        };
992        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
993        p.audit_heads.insert(rec_hash.into(), audit_hash);
994        p.status_index.insert(rec_hash.into(), to);
995        if to == RecStatus::Rejected {
996            if let Ok(rec) = load_rec(sub, rec_hash) {
997                // Exponential backoff keyed on dedup_key: 7d, 14d, 28d, … capped
998                // at 90d, so a finding a reviewer keeps rejecting stops
999                // re-surfacing on a fixed 7d cadence (was a flat 7d despite the
1000                // "doubling" comment).
1001                const BASE_MS: i64 = 7 * 86_400_000;
1002                const CAP_MS: i64 = 90 * 86_400_000;
1003                let strikes = p.cooldown_strikes.entry(rec.dedup_key.clone()).or_insert(0);
1004                let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
1005                *strikes = strikes.saturating_add(1);
1006                p.cooldowns.insert(rec.dedup_key, now_ms + interval);
1007            }
1008        }
1009        sub.store_state(&p.to_value()?)?;
1010        Ok(())
1011    }
1012
1013    /// Check everything [`apply`](Self::apply) would refuse on, without writing
1014    /// anything.
1015    ///
1016    /// This exists for the fused approve-and-apply callers (the bindings'
1017    /// `apply_recommendation`). Recording the approval first and *then* hitting
1018    /// the destructive gate strands the recommendation in `approved`, which has
1019    /// no exit but `applied` or `expired` — `approved → rejected` is not a
1020    /// legal transition — so a refused apply left the reviewer unable to
1021    /// dismiss it. Ask first, then approve.
1022    ///
1023    /// Deliberately does not check the lifecycle transition: the caller is
1024    /// about to make it legal by approving.
1025    /// `has_gating` is whether the caller will supply a gating run at apply:
1026    /// a gated revision (code or adapter) without one is refused HERE, before
1027    /// a fused approve-and-apply records the approval — `approved` has no
1028    /// exit but `applied` or `expired`, so asking after would strand it.
1029    pub fn preflight_apply<S: OmsSubstrate>(
1030        &self,
1031        sub: &S,
1032        rec_hash: &str,
1033        scopes: &ScopeSet,
1034        allow_destructive: bool,
1035        has_gating: bool,
1036    ) -> Result<()> {
1037        if !scopes.has(Scope::Apply) {
1038            return Err(Error::ScopeDenied("apply".into()));
1039        }
1040        let rec = load_rec(sub, rec_hash)?;
1041        if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1042            return Err(Error::DestructiveGated(
1043                "destructive apply requires admin scope + allow_destructive".into(),
1044            ));
1045        }
1046        ensure_executable(rec.action_kind, &rec.proposal)?;
1047        if requires_gating(rec.action_kind) && !has_gating {
1048            return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1049        }
1050        Ok(())
1051    }
1052
1053    /// Apply an approved recommendation. Requires `apply`; destructive payloads
1054    /// additionally require `admin` + `allow_destructive`. Records the applied
1055    /// info (inverse plan) for rollback.
1056    #[allow(clippy::too_many_arguments)]
1057    pub fn apply<S: OmsSubstrate>(
1058        &self,
1059        sub: &mut S,
1060        rec_hash: &str,
1061        actor: &str,
1062        observer: ObserverType,
1063        scopes: &ScopeSet,
1064        because: &str,
1065        allow_destructive: bool,
1066        now_ms: i64,
1067    ) -> Result<AppliedRecord> {
1068        self.apply_inner(
1069            sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1070        )
1071    }
1072
1073    /// Load the gating evidence for `rec_hash` from the RECORDED
1074    /// `mg:eval_run` summary named by `run_id` — the one loader every
1075    /// surface (CLI, bindings, MCP, HTTP) shares, so the stats that admit a
1076    /// gated revision can never come from a caller: they are read back from
1077    /// the Fact `areev eval run` journaled, within the recommendation's own
1078    /// pinned evalset (a run id from a different evalset simply isn't found).
1079    pub fn gating_evidence<S: OmsSubstrate>(
1080        &self,
1081        sub: &S,
1082        rec_hash: &str,
1083        run_id: &str,
1084    ) -> Result<crate::recommendation::GatingEvidence> {
1085        let rec = self
1086            .recommendations(sub, None)?
1087            .into_iter()
1088            .find(|r| r.hash == rec_hash)
1089            .ok_or_else(|| {
1090                Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1091            })?;
1092        let pin = rec.evalset_hash.ok_or_else(|| {
1093            Error::InvalidProposal(
1094                "this recommendation pins no evalset — a gating run applies only \
1095                 to code and adapter revisions"
1096                    .into(),
1097            )
1098        })?;
1099        // Read through the shared evalset reader, so the run that GATES an
1100        // apply and the runs that later JUDGE it are parsed by exactly one
1101        // piece of code — two parsers drifting apart would let a rule be
1102        // admitted on one reading of a summary and measured on another.
1103        match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1104            Some(run) => Ok(crate::recommendation::GatingEvidence {
1105                evalset_hash: pin,
1106                run_id: run.run_id,
1107                passed: run.passed,
1108                failed: run.failed,
1109            }),
1110            None => Err(Error::InvalidProposal(format!(
1111                "no recorded gate run '{run_id}' for evalset {pin} — run \
1112                 `areev eval run --evalset {pin} ...` first"
1113            ))),
1114        }
1115    }
1116
1117    /// Apply WITH the §7.4 evalset-run edge — the only path that can apply a
1118    /// gated (code or adapter) revision. The evidence is validated against
1119    /// the recommendation's pin and recorded on the audit Observation.
1120    #[allow(clippy::too_many_arguments)]
1121    pub fn apply_gated<S: OmsSubstrate>(
1122        &self,
1123        sub: &mut S,
1124        rec_hash: &str,
1125        actor: &str,
1126        observer: ObserverType,
1127        scopes: &ScopeSet,
1128        because: &str,
1129        allow_destructive: bool,
1130        gating: &crate::recommendation::GatingEvidence,
1131        now_ms: i64,
1132    ) -> Result<AppliedRecord> {
1133        self.apply_inner(
1134            sub,
1135            rec_hash,
1136            actor,
1137            observer,
1138            scopes,
1139            because,
1140            allow_destructive,
1141            Some(gating),
1142            now_ms,
1143        )
1144    }
1145
1146    #[allow(clippy::too_many_arguments)]
1147    fn apply_inner<S: OmsSubstrate>(
1148        &self,
1149        sub: &mut S,
1150        rec_hash: &str,
1151        actor: &str,
1152        observer: ObserverType,
1153        scopes: &ScopeSet,
1154        because: &str,
1155        allow_destructive: bool,
1156        gating: Option<&crate::recommendation::GatingEvidence>,
1157        now_ms: i64,
1158    ) -> Result<AppliedRecord> {
1159        if !scopes.has(Scope::Apply) {
1160            return Err(Error::ScopeDenied("apply".into()));
1161        }
1162        let because = validate_because(because)?;
1163        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1164        let status = *p
1165            .status_index
1166            .get(rec_hash)
1167            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1168        if !status.can_transition_to(RecStatus::Applied, false) {
1169            return Err(Error::LifecycleViolation(format!(
1170                "{} -> applied (approve first)",
1171                status.as_str()
1172            )));
1173        }
1174        let rec = load_rec(sub, rec_hash)?;
1175        if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1176            return Err(Error::DestructiveGated(
1177                "destructive apply requires admin scope + allow_destructive".into(),
1178            ));
1179        }
1180        // §7.4: a code or adapter revision applies ONLY through the
1181        // evalset-run edge — the pin must match, the pinned evalset must
1182        // still be LIVE (a superseded evalset invalidates in-flight
1183        // recommendations: they re-gate), and a failing gate admits nothing.
1184        if requires_gating(rec.action_kind) {
1185            let g = gating
1186                .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1187            let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1188            if g.evalset_hash != pin {
1189                return Err(Error::InvalidProposal(format!(
1190                    "gating ran evalset {} but the recommendation is pinned \
1191                     to {pin} (Rule E1)",
1192                    g.evalset_hash
1193                )));
1194            }
1195            match sub.grain(pin)? {
1196                Some(evalset) if evalset.is_live() => {}
1197                Some(_) => {
1198                    return Err(Error::InvalidProposal(
1199                        "the pinned evalset was superseded after gating — \
1200                         the recommendation must re-gate (Rule E1)"
1201                            .into(),
1202                    ))
1203                }
1204                None => {
1205                    return Err(Error::InvalidProposal(format!(
1206                        "pinned evalset {pin} not found in the substrate"
1207                    )))
1208                }
1209            }
1210            if g.failed > 0 {
1211                return Err(Error::InvalidProposal(format!(
1212                    "the gating run failed {}/{} cases — a failing gate \
1213                     admits nothing",
1214                    g.failed,
1215                    g.passed + g.failed
1216                )));
1217            }
1218        }
1219
1220        // Execute the proposal.
1221        let mut created = Vec::new();
1222        // The inverse of a change that creates no grain — see
1223        // `AppliedRecord::inverse_cal`. Captured BEFORE execution, because
1224        // afterwards the previous definition is gone.
1225        let mut inverse_cal: Option<String> = None;
1226        match &rec.proposal {
1227            Proposal::Cal { cal } => {
1228                for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1229                    if !is_definition_statement(line) {
1230                        continue;
1231                    }
1232                    match sub.definition_inverse(line)? {
1233                        Some(inv) => inverse_cal = Some(inv),
1234                        None => {
1235                            return Err(Error::InvalidProposal(format!(
1236                                "this substrate cannot record a rollback inverse for {line:?}; \
1237                                 a definition rewrite that ROLLBACK could not undo is refused \
1238                                 rather than applied"
1239                            )))
1240                        }
1241                    }
1242                }
1243                let rows = sub.execute_cal(cal)?;
1244                for r in rows {
1245                    if let Some(h) = r.get("hash").and_then(Value::as_str) {
1246                        created.push(h.to_string());
1247                    }
1248                }
1249            }
1250            // The engine has no executable Edit primitive. Marking this
1251            // Applied used to be a lie (and rollback had no inverse).
1252            Proposal::Edit { .. } => {
1253                return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1254            }
1255            // A gated code or adapter revision EXECUTES by writing the
1256            // promotion grain: an immutable record that this target now
1257            // resolves to the proposed code (a §7.4 blob address) or adapter
1258            // (the tuning seam's registry tuple). Hosts read the promotion
1259            // to re-resolve — `(tool:X, mg:code_promotion)` or
1260            // `(model:X, mg:adapter_promotion)`; retracting it is the
1261            // rollback inverse, so the apply is rollbackable end-to-end.
1262            Proposal::Data { data } if requires_gating(rec.action_kind) => {
1263                let relation = if rec.action_kind == ActionKind::AdapterRevision {
1264                    "mg:adapter_promotion"
1265                } else {
1266                    "mg:code_promotion"
1267                };
1268                let mut spec = crate::substrate::GrainSpec::new(
1269                    crate::model::grain_type::FACT,
1270                    LOOP_NS,
1271                )
1272                .with_field("subject", rec.target_ref.clone())
1273                .with_field("relation", relation)
1274                .with_field(
1275                    "object",
1276                    serde_json::to_string(data)
1277                        .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1278                )
1279                .with_field("rec_hash", rec_hash.to_string());
1280                if let Some(g) = gating {
1281                    spec = spec
1282                        .with_field("gating_evalset", g.evalset_hash.clone())
1283                        .with_field("gating_run_id", g.run_id.clone());
1284                }
1285                created.push(sub.put_grain(&spec)?);
1286            }
1287            Proposal::Data { data } => {
1288                // OutcomeReview is the one executable Data shape: its
1289                // `revert_of` points at an earlier applied recommendation.
1290                // Reuse the ordinary rollback path so the created hashes are
1291                // really retracted and the original lifecycle/audit advances.
1292                let revert_of = data
1293                    .get("revert_of")
1294                    .and_then(Value::as_str)
1295                    .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1296                self.rollback(
1297                    sub,
1298                    revert_of,
1299                    actor,
1300                    observer,
1301                    scopes,
1302                    &because,
1303                    now_ms,
1304                )?;
1305                // rollback stored a newer lifecycle state; merge this apply
1306                // into that state rather than overwriting the rollback.
1307                p = LoopPersisted::from_value(sub.load_state()?)?;
1308            }
1309        }
1310
1311        let applied = AppliedRecord {
1312            applied_at_ms: now_ms,
1313            target_ref: rec.target_ref.clone(),
1314            rollbackable: rec.rollbackable,
1315            created_hashes: created,
1316            inverse_cal,
1317            metric: rec.metric.clone(),
1318        };
1319        let prev = p.audit_heads.get(rec_hash).cloned();
1320        let audit = AuditRecord {
1321            rec_hash: rec_hash.into(),
1322            from: Some(status),
1323            to: RecStatus::Applied,
1324            actor: actor.into(),
1325            observer_type: observer,
1326            because,
1327            previous_audit_hash: prev,
1328            gating: gating.cloned(),
1329            at_ms: now_ms,
1330        };
1331        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1332        p.audit_heads.insert(rec_hash.into(), audit_hash);
1333        p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1334        p.applied.insert(rec_hash.into(), applied.clone());
1335        sub.store_state(&p.to_value()?)?;
1336        Ok(applied)
1337    }
1338
1339    /// Roll back an applied recommendation by retracting the grains it created.
1340    /// Fails for non-rollbackable applies (e.g. FORGET).
1341    #[allow(clippy::too_many_arguments)]
1342    pub fn rollback<S: OmsSubstrate>(
1343        &self,
1344        sub: &mut S,
1345        rec_hash: &str,
1346        actor: &str,
1347        observer: ObserverType,
1348        scopes: &ScopeSet,
1349        because: &str,
1350        now_ms: i64,
1351    ) -> Result<()> {
1352        if !scopes.has(Scope::Apply) {
1353            return Err(Error::ScopeDenied("apply".into()));
1354        }
1355        let because = validate_because(because)?;
1356        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1357        let status = *p
1358            .status_index
1359            .get(rec_hash)
1360            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1361        if !status.can_transition_to(RecStatus::RolledBack, false) {
1362            return Err(Error::LifecycleViolation(format!(
1363                "{} -> rolled_back",
1364                status.as_str()
1365            )));
1366        }
1367        let applied = p
1368            .applied
1369            .get(rec_hash)
1370            .cloned()
1371            .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1372        if !applied.rollbackable {
1373            return Err(Error::LifecycleViolation(
1374                "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1375            ));
1376        }
1377        for h in &applied.created_hashes {
1378            sub.retract(h, &format!("rollback of {rec_hash}"))?;
1379        }
1380        // A definition rewrite creates no grain, so retracting `created_hashes`
1381        // undoes nothing. Restoring it means re-running the statement captured
1382        // at apply time — the previous definition, or a DROP when there was
1383        // none. Runs BEFORE the audit is written, so a failed restore leaves
1384        // the recommendation `applied` (still true) rather than recording a
1385        // rollback that did not happen.
1386        if let Some(inverse) = &applied.inverse_cal {
1387            sub.execute_cal(inverse)?;
1388        }
1389        let prev = p.audit_heads.get(rec_hash).cloned();
1390        let audit = AuditRecord {
1391            rec_hash: rec_hash.into(),
1392            from: Some(status),
1393            to: RecStatus::RolledBack,
1394            actor: actor.into(),
1395            observer_type: observer,
1396            because,
1397            previous_audit_hash: prev,
1398            gating: None,
1399            at_ms: now_ms,
1400        };
1401        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1402        p.audit_heads.insert(rec_hash.into(), audit_hash);
1403        p.status_index
1404            .insert(rec_hash.into(), RecStatus::RolledBack);
1405        sub.store_state(&p.to_value()?)?;
1406        Ok(())
1407    }
1408
1409    /// List stored recommendations, optionally filtered by status. Status comes
1410    /// from the rebuildable index, not the immutable grain body. Ordered for
1411    /// review triage — highest severity first, then oldest first — and stable
1412    /// across runs for identical input.
1413    pub fn recommendations<S: OmsSubstrate>(
1414        &self,
1415        sub: &S,
1416        status_filter: Option<RecStatus>,
1417    ) -> Result<Vec<Recommendation>> {
1418        let p = LoopPersisted::from_value(sub.load_state()?)?;
1419        let grains = sub.grains_of_type(
1420            crate::model::grain_type::RECOMMENDATION,
1421            Some(LOOP_NS),
1422            ReadOpts {
1423                live_only: false,
1424                since_ms: None,
1425            },
1426        )?;
1427        let mut out = Vec::new();
1428        for g in grains {
1429            let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1430            rec.status = p
1431                .status_index
1432                .get(&g.hash)
1433                .copied()
1434                .unwrap_or(RecStatus::Pending);
1435            if let Some(f) = status_filter {
1436                if rec.status != f {
1437                    continue;
1438                }
1439            }
1440            out.push(rec);
1441        }
1442        // Review-queue order: worst first, then oldest first. Hash is only the
1443        // final tiebreak — sorting by it alone is deterministic per run but
1444        // meaningless across runs, because a grain's hash covers its timestamp,
1445        // so an identical queue comes back in a different order every time.
1446        // `dedup_key` is the last tiebreak that actually decides anything: it is
1447        // content-derived and stable across runs, whereas findings proposed in
1448        // the same sweep routinely share a `created_at_ms`. Hash trails it only
1449        // to make the ordering total.
1450        out.sort_by(|a, b| {
1451            b.severity
1452                .cmp(&a.severity)
1453                .then(a.created_at_ms.cmp(&b.created_at_ms))
1454                .then(a.dedup_key.cmp(&b.dedup_key))
1455                .then(a.hash.cmp(&b.hash))
1456        });
1457        Ok(out)
1458    }
1459
1460    /// Per-analyzer effective settings for the Setup view: the manifest facts
1461    /// merged with the file-config (override or manifest default). Read-only.
1462    pub fn analyzer_settings<S: OmsSubstrate>(
1463        &self,
1464        sub: &S,
1465    ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1466        let p = LoopPersisted::from_value(sub.load_state()?)?;
1467        Ok(self
1468            .analyzers
1469            .iter()
1470            .map(|a| {
1471                let m = a.manifest();
1472                let cfg = p.config.get(&m.id);
1473                crate::config::AnalyzerSetting {
1474                    id: m.id.clone(),
1475                    title: m.title.clone(),
1476                    description: m.description.clone(),
1477                    tier: format!("{:?}", m.tier),
1478                    trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1479                    default_on: m.default_on,
1480                    enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1481                    severity_floor: cfg
1482                        .and_then(|c| c.severity_floor)
1483                        .map(|s| s.as_str().to_string()),
1484                }
1485            })
1486            .collect())
1487    }
1488
1489    /// Update one analyzer's file-config (enable/disable, severity floor, param
1490    /// overrides, namespace scoping). Requires `Admin`. Params are validated
1491    /// against the analyzer's manifest first (unknown keys rejected, fail-closed),
1492    /// and the analyzer must exist. Returns the merged config as stored. This is
1493    /// the only write into `persisted.config` — the config layer, never a grain.
1494    pub fn set_analyzer_config<S: OmsSubstrate>(
1495        &self,
1496        sub: &mut S,
1497        analyzer_id: &str,
1498        update: crate::config::AnalyzerConfigUpdate,
1499        scopes: &ScopeSet,
1500    ) -> Result<crate::config::AnalyzerConfig> {
1501        if !scopes.has(Scope::Admin) {
1502            return Err(Error::ScopeDenied("admin".into()));
1503        }
1504        let manifest = self
1505            .analyzers
1506            .iter()
1507            .map(|a| a.manifest())
1508            .find(|m| m.id == analyzer_id)
1509            .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1510        // Validate params against the manifest BEFORE touching state.
1511        if let Some(params) = &update.params {
1512            manifest.resolve_params(params)?;
1513        }
1514        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1515        let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1516        if let Some(enabled) = update.enabled {
1517            cfg.enabled = Some(enabled);
1518        }
1519        if update.clear_floor {
1520            cfg.severity_floor = None;
1521        } else if let Some(floor) = update.severity_floor {
1522            cfg.severity_floor = Some(floor);
1523        }
1524        if let Some(params) = update.params {
1525            cfg.params = params;
1526        }
1527        if let Some(ns) = update.namespaces {
1528            cfg.namespaces = ns;
1529        }
1530        let stored = cfg.clone();
1531        sub.store_state(&p.to_value()?)?;
1532        Ok(stored)
1533    }
1534
1535    /// The measured outcome time series (the Verify gate's history) across all
1536    /// recommendations, ordered by when each checkpoint was measured.
1537    pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1538        let p = LoopPersisted::from_value(sub.load_state()?)?;
1539        let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1540        // `metric` and `rec_hash` break the tie: checkpoints measured in the
1541        // same sweep share a `measured_at_ms`, and without a tiebreak the order
1542        // falls through to the map's rec_hash ordering, which shifts every run.
1543        out.sort_by(|a, b| {
1544            a.measured_at_ms
1545                .cmp(&b.measured_at_ms)
1546                .then(a.horizon_ms.cmp(&b.horizon_ms))
1547                .then(a.metric.cmp(&b.metric))
1548                .then(a.rec_hash.cmp(&b.rec_hash))
1549        });
1550        Ok(out)
1551    }
1552
1553    /// A health snapshot — when the loop last ran, how much is un-analyzed
1554    /// since, and the queue counts. Lets a host surface "the loop may be stale"
1555    /// so a forgotten SessionEnd hook / cron doesn't silently kill it.
1556    pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1557        let p = LoopPersisted::from_value(sub.load_state()?)?;
1558        let (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
1559        let recs = self.recommendations(sub, None)?;
1560        let mut pending = 0;
1561        let mut applied = 0;
1562        for r in &recs {
1563            match r.status {
1564                RecStatus::Pending => pending += 1,
1565                RecStatus::Applied => applied += 1,
1566                _ => {}
1567            }
1568        }
1569        // Stale if it has never run, or it's been a while / a lot has piled up.
1570        let stale = match p.state.last_run_ms {
1571            None => true,
1572            Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1573        };
1574        Ok(Health {
1575            last_run_ms: p.state.last_run_ms,
1576            grains_since_run,
1577            error_events_since_run,
1578            pending,
1579            applied,
1580            total: recs.len() as u64,
1581            stale,
1582        })
1583    }
1584
1585    /// Approval-rate metric for `origin = llm` recommendations (reflection
1586    /// design §6b) — the live field-quality signal that accrues off the audit
1587    /// chain: what fraction of the model's *surfaced* proposals a reviewer
1588    /// accepts. Complements the offline Effective-Reliability eval.
1589    pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1590        let recs = self.recommendations(sub, None)?;
1591        let mut m = LlmMetrics::default();
1592        for r in &recs {
1593            if !matches!(r.origin, Origin::Llm { .. }) {
1594                continue;
1595            }
1596            m.proposed += 1;
1597            match r.status {
1598                RecStatus::Pending => m.pending += 1,
1599                RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1600                RecStatus::Rejected => m.rejected += 1,
1601                RecStatus::Expired => {}
1602            }
1603        }
1604        let decided = m.approved + m.rejected;
1605        m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1606        Ok(m)
1607    }
1608}
1609
1610/// A health snapshot for the backend's self-improvement loop.
1611#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1612pub struct Health {
1613    #[serde(skip_serializing_if = "Option::is_none")]
1614    pub last_run_ms: Option<i64>,
1615    pub grains_since_run: u64,
1616    pub error_events_since_run: u64,
1617    pub pending: u64,
1618    pub applied: u64,
1619    pub total: u64,
1620    /// True when the loop looks stalled (never run, or ≥7d / ≥100 new grains
1621    /// since the last run) — a nudge that a trigger may be unwired.
1622    pub stale: bool,
1623}
1624
1625/// Approval-rate metric for `origin = llm` recommendations (reflection §6b).
1626#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1627pub struct LlmMetrics {
1628    /// Total llm-origin recommendations ever stored (those that survived the
1629    /// verifier and reached the queue).
1630    pub proposed: u64,
1631    pub pending: u64,
1632    /// Approved + Applied + RolledBack (a reviewer said yes at least once).
1633    pub approved: u64,
1634    pub rejected: u64,
1635    /// approved / (approved + rejected); `None` until at least one is decided.
1636    #[serde(skip_serializing_if = "Option::is_none")]
1637    pub approval_rate: Option<f64>,
1638}
1639
1640/// Re-measure applied recommendations at each **checkpoint** past due, via the
1641/// engine's typed reads — no CAL-scalar round-trip. A recommendation
1642/// accumulates one `OutcomeResult` per horizon (measured once each), forming a
1643/// time series, so a late regression (held at 1d, regressed at 30d) is caught.
1644/// Only *regressed* checkpoints feed the outcome analyzer (→ a revert).
1645/// Unknown metric kinds are skipped, never faked.
1646fn measure_outcomes<S: OmsSubstrate>(
1647    sub: &S,
1648    p: &mut LoopPersisted,
1649    now_ms: i64,
1650) -> Result<Vec<OutcomeInput>> {
1651    // Collect all due (recommendation, horizon) checkpoints first.
1652    let mut due: Vec<(String, crate::config::AppliedRecord, i64)> = Vec::new();
1653    for (h, a) in &p.applied {
1654        if p.status_index.get(h) != Some(&RecStatus::Applied) {
1655            continue;
1656        }
1657        let Some(metric) = &a.metric else { continue };
1658        let done = p.measured.get(h).cloned().unwrap_or_default();
1659        for horizon in metric.horizons() {
1660            if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
1661                due.push((h.clone(), a.clone(), horizon));
1662            }
1663        }
1664    }
1665
1666    let mut out = Vec::new();
1667    for (rec_hash, applied, horizon) in due {
1668        let metric = applied.metric.as_ref().unwrap();
1669        let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1670            continue; // metric kind not yet re-measurable
1671        };
1672        let regressed = crate::recommendation::is_regression(
1673            metric.baseline,
1674            current,
1675            metric.higher_is_better,
1676        );
1677        p.outcomes.entry(rec_hash.clone()).or_default().push(
1678            crate::recommendation::OutcomeResult {
1679                rec_hash: rec_hash.clone(),
1680                metric: metric.metric.clone(),
1681                baseline: metric.baseline,
1682                current,
1683                verdict: if regressed { "regressed" } else { "held" }.into(),
1684                horizon_ms: horizon,
1685                measured_at_ms: now_ms,
1686            },
1687        );
1688        p.measured.entry(rec_hash.clone()).or_default().push(horizon);
1689        if regressed {
1690            out.push(OutcomeInput {
1691                rec_hash,
1692                target_ref: applied.target_ref.clone(),
1693                metric: metric.metric.clone(),
1694                baseline: metric.baseline,
1695                current,
1696                unit: metric.unit.clone(),
1697                higher_is_better: metric.higher_is_better,
1698            });
1699        }
1700    }
1701    Ok(out)
1702}
1703
1704/// Typed re-measurement for the fixed set of metric kinds the engine knows.
1705pub(crate) fn measure_metric<S: SubstrateRead>(
1706    sub: &S,
1707    metric: &crate::recommendation::MetricSnapshot,
1708    since_ms: i64,
1709) -> Result<Option<f64>> {
1710    match metric.metric.as_str() {
1711        // How many times did this tool fail again *with the same signature*
1712        // after the lesson was applied? Scoped to the signature (metric.relation)
1713        // so an unrelated later failure of the same tool is not read as a
1714        // regression of this specific lesson.
1715        "tool_error_recurrence" => {
1716            let Some(tool) = &metric.subject else { return Ok(None) };
1717            let tools = sub.grains_of_type(
1718                crate::model::grain_type::TOOL,
1719                None,
1720                ReadOpts { live_only: true, since_ms: Some(since_ms) },
1721            )?;
1722            let n = tools
1723                .iter()
1724                .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
1725                .filter(|t| {
1726                    // No stored signature (legacy metric) → fall back to the
1727                    // whole-tool count so old recommendations still measure.
1728                    metric.relation.as_deref().is_none_or(|sig| {
1729                        crate::analyzers::tool_failure::normalize_signature(
1730                            t.tool_content().unwrap_or(""),
1731                        ) == sig
1732                    })
1733                })
1734                .count();
1735            Ok(Some(n as f64))
1736        }
1737        // After a resolve-to-latest, does the subject again hold more than one
1738        // live value under the functional relation? Live-state read (no since
1739        // filter): the excess beyond one distinct object is the regression.
1740        "contradiction_recurrence" => {
1741            let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
1742                return Ok(None);
1743            };
1744            let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
1745            let distinct: BTreeSet<String> = facts
1746                .iter()
1747                .filter(|f| {
1748                    f.fact_relation()
1749                        .is_some_and(|r| normalize_ident(r) == *relation)
1750                })
1751                .filter_map(|f| f.fact_object().map(normalize_ident))
1752                .collect();
1753            Ok(Some(distinct.len().saturating_sub(1) as f64))
1754        }
1755        // `evalset:<hash>:<field>` — external correctness, measured the only
1756        // way that keeps the honesty rule ("Areev Loop improves the agent's
1757        // memory, not its outputs") intact: an evalset run is an INTERNAL,
1758        // BOUNDED, ATTRIBUTABLE measurement. The engine never runs the evalset;
1759        // it only reads summaries a host journaled with `areev eval run`.
1760        //
1761        // `since_ms` is the apply time, and it is load-bearing rather than an
1762        // optimization: a run journaled BEFORE the apply cannot be evidence of
1763        // what applying did. Comparing the baseline run against itself would
1764        // report "held" forever — a fabricated receipt, which is worse than no
1765        // receipt at all. No run since the apply → `None` → not yet
1766        // measurable, and the checkpoint stays due.
1767        m if m.starts_with("evalset:") => {
1768            let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
1769                return Ok(None);
1770            };
1771            let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
1772                return Ok(None);
1773            };
1774            // `failed`/`passed`/`total` are promoted so a metric can be written
1775            // against any evalset without the host having to add fields;
1776            // anything else is read from the summary the host did write.
1777            Ok(match field {
1778                "failed" => Some(run.failed as f64),
1779                "passed" => Some(run.passed as f64),
1780                "total" => Some(run.total() as f64),
1781                "error_rate" => match run.total() {
1782                    0 => None, // no cases ran: undefined, not zero
1783                    t => Some(run.failed as f64 / t as f64),
1784                },
1785                other => run.field(other),
1786            })
1787        }
1788        _ => Ok(None),
1789    }
1790}
1791
1792/// Live facts for one normalized (namespace?, subject) — the shared scope of
1793/// the fact-shaped recurrence metrics. `namespace: None` spans all namespaces.
1794fn scoped_live_facts<S: SubstrateRead>(
1795    sub: &S,
1796    namespace: Option<&str>,
1797    subject: &str,
1798) -> Result<Vec<GrainRecord>> {
1799    let facts = sub.grains_of_type(
1800        crate::model::grain_type::FACT,
1801        None,
1802        ReadOpts { live_only: true, since_ms: None },
1803    )?;
1804    Ok(facts
1805        .into_iter()
1806        .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
1807        .filter(|f| {
1808            f.fact_subject()
1809                .is_some_and(|s| normalize_ident(s) == subject)
1810        })
1811        .collect())
1812}
1813
1814// --- free helpers ---
1815
1816/// The §6.3 exact-equality check: a SUPERSEDE is *value-identical* when every
1817/// replacement field equals the superseded grain's value — strings after
1818/// case-fold/trim (upstream NFC is an OMS invariant), `namespace` against the
1819/// grain's own namespace, everything else exactly. This is what makes an
1820/// auto-applied consolidation provably information-preserving; a
1821/// near-duplicate (an observation body off by one token) fails it and stays
1822/// pending for human review. Fails closed: an unrecognized line shape, an
1823/// empty replacement, a missing grain, or a field the original never had all
1824/// disqualify.
1825fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
1826    let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
1827        return false;
1828    };
1829    if fields.is_empty() {
1830        return false;
1831    }
1832    let Ok(Some(grain)) = sub.grain(&target) else {
1833        return false;
1834    };
1835    // An expiry lives OUTSIDE `fields`, so the replacement (built from a fixed
1836    // field set) can never carry it — consolidating away a grain that has a
1837    // valid_to would silently drop the expiry, invisibly to the field-by-field
1838    // check below. Fail closed. (A dup that additionally carries an extra
1839    // *content* field the replacement omits is a real but subtler info-loss;
1840    // catching it soundly needs a full canonical-vs-extra comparison rather
1841    // than this replacement-scoped check, since a real fact's fields also carry
1842    // OMS metadata like `confidence` that consolidation legitimately keeps —
1843    // left as a follow-up so this narrow fix can't block valid consolidations.)
1844    if grain.valid_to_ms.is_some() {
1845        return false;
1846    }
1847    // Forward check: every replacement field equals the grain's value.
1848    fields.iter().all(|(k, v)| {
1849        if k == "namespace" {
1850            return v
1851                .as_str()
1852                .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
1853        }
1854        match (v, grain.fields.get(k)) {
1855            (Value::String(a), Some(Value::String(b))) => {
1856                normalize_ident(a) == normalize_ident(b)
1857            }
1858            (a, Some(b)) => a == b,
1859            (_, None) => false,
1860        }
1861    })
1862}
1863
1864/// The confidence floor (§5.4): a verified draft below this is dropped. The
1865/// verifier's calibrated confidence is the gate, not the proposer's self-report.
1866const MIN_LLM_CONFIDENCE: f64 = 0.75;
1867
1868/// The fixed DISCOVER instruction (§5.1). The scoring rule makes "nothing to
1869/// report" a first-class, zero-penalty answer — the structural antidote to
1870/// over-generation. Kept in its own request field so it never interleaves with
1871/// (attacker-influenced) evidence text.
1872const DISCOVER_INSTRUCTIONS: &str = "You review an agent's memory for quality. \
1873Given deterministic findings and the evidence they cite, propose ADDITIONAL \
1874findings the deterministic checks would miss (e.g. a semantic contradiction, a \
1875stale assumption, a duplicated meaning). SCORING: propose a finding ONLY if you \
1876are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
1877useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
1878earns 0. When in doubt, propose nothing — an empty list is the correct answer \
1879when there is nothing worth flagging. The 'approved' and 'rejected' lists, when \
1880present, show findings this reviewer recently accepted or rejected — prefer the \
1881kind they accept and avoid the kind they reject. Every proposal MUST cite one \
1882or more evidence hashes from the bundle, target a memory entity, and include \
1883your confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
1884\"target\":\"entity:<ns>/<subject>\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
1885\"confidence\":0.0}]}. Propose nothing you cannot ground in the evidence.";
1886
1887/// The fixed GROUND instruction (§5.2): verify the finding's factual PREMISES
1888/// are real (anti-fabrication), while allowing an inference. A self-improvement
1889/// finding may reason BEYOND the evidence (e.g. 'HQ=SF and country=Germany are
1890/// inconsistent'); grounding checks the premises (HQ=SF, country=Germany) are in
1891/// the evidence — the soundness of the inference is VERIFY's job, not this one.
1892const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
1893fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
1894the facts it relies on are actually present in the cited evidence, NOT that its \
1895conclusion is stated verbatim. Decompose the finding into the factual claims it \
1896depends on. Mark supported=true when those facts are present in the evidence \
1897(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
1898on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
1899different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
1900
1901/// The fixed VERIFY instruction (§5.3): adversarial, abstention-biased.
1902const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
1903each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
1904never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
1905SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
1906'possible' findings with no concrete defect, and reject any claimed \
1907inconsistency or contradiction that is not backed by at least two actually \
1908conflicting facts in the cited evidence. (2) Context — does the finding \
1909correctly read its cited evidence, or misinterpret what the grains say? Keep a \
1910finding when it names a genuine, specific problem grounded in its evidence and \
1911materially useful to a human reviewer; otherwise reject it, and default to \
1912keep=false when uncertain. Do NOT reject a finding for being 'already known', \
1913redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
1914grounded cross-fact inconsistency with two conflicting facts is exactly what to \
1915KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
1916{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
1917
1918/// The fixed ENRICH instruction.
1919const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
1920guidance note to help a human reviewer decide. Do not restate the finding. Return \
1921JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
1922
1923/// Add a grain to the evidence bundle — deduplicated, bounded at 64, text
1924/// capped. Shared by the deterministic-citation and recent-grain seeding.
1925fn push_evidence(
1926    evidence: &mut Vec<crate::llm::EvidenceItem>,
1927    bundle: &mut BTreeSet<String>,
1928    g: &GrainRecord,
1929) {
1930    if evidence.len() < 64 && bundle.insert(g.hash.clone()) {
1931        evidence.push(crate::llm::EvidenceItem {
1932            hash: g.hash.clone(),
1933            grain_type: g.grain_type.clone(),
1934            text: crate::llm::cap(&grain_brief(g), 400),
1935        });
1936    }
1937}
1938
1939/// A short human-readable projection of a grain for the evidence bundle.
1940fn grain_brief(g: &GrainRecord) -> String {
1941    if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
1942        return format!("{s} {r} {o}");
1943    }
1944    for key in ["content", "body", "text", "summary"] {
1945        if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
1946            if !v.is_empty() {
1947                return v.to_string();
1948            }
1949        }
1950    }
1951    String::new()
1952}
1953
1954/// Stamp a validated DISCOVER draft as an `origin = llm` recommendation. LLM
1955/// drafts are always advisory `Flag`s carrying `Proposal::Data` (never an
1956/// executable CAL mutation), lower-confidence, and — via `Origin::Llm` and a
1957/// no-manifest analyzer id — structurally ineligible for auto-apply.
1958fn stamp_llm(
1959    model: &str,
1960    d: &crate::llm::LlmDraft,
1961    target_ref: String,
1962    cited: Vec<String>,
1963    confidence: f64,
1964    now_ms: i64,
1965) -> Recommendation {
1966    let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
1967    let mut args = serde_json::Map::new();
1968    args.insert("text".into(), Value::from(summary_text));
1969    let guidance = if d.guidance.trim().is_empty() {
1970        None
1971    } else {
1972        Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
1973    };
1974    let action = ActionKind::Flag;
1975    let mut data = serde_json::Map::new();
1976    data.insert("source".into(), Value::from("llm"));
1977    Recommendation {
1978        hash: String::new(),
1979        analyzer: "loop.llm/1".to_string(),
1980        params_snapshot: serde_json::Map::new(),
1981        origin: Origin::Llm { model: model.to_string() },
1982        target_ref: target_ref.clone(),
1983        action_kind: action,
1984        dedup_key: dedup_key("llm", &target_ref, action),
1985        summary: Summary::new("llm.discover", args),
1986        severity: Severity::Low,
1987        proposal: Proposal::Data { data },
1988        destructive: false,
1989        rollbackable: false,
1990        evidence: cited,
1991        evidence_query: None,
1992        metric: None,
1993        // The verifier's calibrated confidence — not a hardcoded default.
1994        confidence: confidence.clamp(0.0, 1.0),
1995        importance: 0.3,
1996        created_at_ms: now_ms,
1997        guidance,
1998        evalset_hash: None,
1999        status: RecStatus::Pending,
2000    }
2001}
2002
2003/// The action kinds that apply ONLY through the evalset-run gating edge
2004/// (§7.4 for tool code; the tuning seam's adapter promotion inherits the
2005/// same rule). One predicate so the gate, the rollbackable stamp, and the
2006/// promotion write can never disagree on membership.
2007fn requires_gating(kind: ActionKind) -> bool {
2008    matches!(
2009        kind,
2010        ActionKind::CodeRevision | ActionKind::AdapterRevision
2011    )
2012}
2013
2014fn stamp(
2015    m: &AnalyzerManifest,
2016    params: &crate::manifest::Params,
2017    d: crate::recommendation::RecDraft,
2018    now_ms: i64,
2019) -> Result<Recommendation> {
2020    let target = TargetRef::parse(&d.target_ref)?;
2021    // Rule E1 at the door (§7.4): an unpinned code revision, a pinned
2022    // non-code target, or a code action on a non-tool target never becomes
2023    // a recommendation at all.
2024    crate::recommendation::validate_code_rules(
2025        d.action_kind,
2026        target.target_class(),
2027        d.evalset_hash.as_deref(),
2028    )?;
2029    let dedup = dedup_key(m.family(), &d.target_ref, d.action_kind);
2030    let destructive = match &d.proposal {
2031        Proposal::Cal { cal } => cal::contains_destructive(cal),
2032        _ => false,
2033    };
2034    let rollbackable = match &d.proposal {
2035        Proposal::Cal { .. } => !destructive,
2036        Proposal::Edit { .. } => false,
2037        // A code or adapter revision applies by WRITING the promotion
2038        // grain; retracting it is the exact inverse — rollbackable by
2039        // construction.
2040        Proposal::Data { .. } => requires_gating(d.action_kind),
2041    };
2042    let mut evidence = d.evidence;
2043    evidence.truncate(MAX_EVIDENCE);
2044    // Provenance follows the analyzer's trust class: a subprocess
2045    // (`--analyzer-cmd`) finding is stamped `Command` — structurally
2046    // auto-apply-ineligible and badged [external] on the recall surface —
2047    // not `Builtin`.
2048    let origin = match m.trust_class {
2049        crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
2050        _ => Origin::Builtin,
2051    };
2052    Ok(Recommendation {
2053        hash: String::new(),
2054        analyzer: m.id.clone(),
2055        params_snapshot: params.snapshot(),
2056        origin,
2057        target_ref: target.as_string(),
2058        action_kind: d.action_kind,
2059        dedup_key: dedup,
2060        summary: d.summary,
2061        severity: d.severity,
2062        proposal: d.proposal,
2063        destructive,
2064        rollbackable,
2065        evidence,
2066        evidence_query: d.evidence_query,
2067        metric: d.metric,
2068        confidence: d.confidence,
2069        importance: d.importance,
2070        created_at_ms: now_ms,
2071        guidance: None,
2072        evalset_hash: d.evalset_hash,
2073        status: RecStatus::Pending,
2074    })
2075}
2076
2077fn validate_because(because: &str) -> Result<String> {
2078    let trimmed = because.trim();
2079    if trimmed.is_empty() {
2080        return Err(Error::InvalidProposal(
2081            "a BECAUSE reason is required".into(),
2082        ));
2083    }
2084    if trimmed.chars().count() > MAX_BECAUSE {
2085        return Err(Error::InvalidProposal(format!(
2086            "BECAUSE exceeds {MAX_BECAUSE} chars"
2087        )));
2088    }
2089    Ok(trimmed.to_string())
2090}
2091
2092fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
2093    for req in &m.requires {
2094        match req {
2095            Capability::Forks if !caps.forks => return Some("forks"),
2096            Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
2097            Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
2098            _ => {}
2099        }
2100    }
2101    None
2102}
2103
2104fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
2105    p.config.get(analyzer_id).and_then(|c| c.severity_floor)
2106}
2107
2108fn gate(
2109    opts: &RunOptions,
2110    p: &LoopPersisted,
2111    new_grains: u64,
2112    new_errors: u64,
2113    now_ms: i64,
2114) -> Option<SkipReason> {
2115    let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
2116    if !any {
2117        return None;
2118    }
2119    let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
2120    let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
2121    let stale_ok = opts
2122        .if_stale_ms
2123        .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
2124    if min_new_ok || min_err_ok || stale_ok {
2125        return None;
2126    }
2127    // Choose the honest reason: staleness-only gate → not_stale, else min_new.
2128    if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
2129        Some(SkipReason::NotStale)
2130    } else {
2131        Some(SkipReason::MinNewNotMet)
2132    }
2133}
2134
2135fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
2136    let opts = ReadOpts {
2137        live_only: false,
2138        since_ms: watermark.map(|w| w + 1),
2139    };
2140    let mut new_grains = 0u64;
2141    let mut new_errors = 0u64;
2142    for t in [
2143        crate::model::grain_type::FACT,
2144        crate::model::grain_type::EVENT,
2145        crate::model::grain_type::TOOL,
2146        crate::model::grain_type::OBSERVATION,
2147    ] {
2148        let g = sub.grains_of_type(t, None, opts)?;
2149        new_grains += g.len() as u64;
2150        // The error gate (--min-new-errors) watches captured tool failures.
2151        if t == crate::model::grain_type::TOOL {
2152            new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
2153        }
2154    }
2155    Ok((new_grains, new_errors))
2156}
2157
2158fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
2159    let grains = sub.grains_of_type(
2160        crate::model::grain_type::RECOMMENDATION,
2161        Some(LOOP_NS),
2162        ReadOpts {
2163            live_only: false,
2164            since_ms: None,
2165        },
2166    )?;
2167    let mut set = BTreeSet::new();
2168    for g in grains {
2169        let status = p
2170            .status_index
2171            .get(&g.hash)
2172            .copied()
2173            .unwrap_or(RecStatus::Pending);
2174        // Pending/approved (still open) and applied (already handled)
2175        // recommendations suppress re-proposal of the same finding. Rejected
2176        // is handled by cooldowns; rolled_back/expired may legitimately
2177        // re-propose (the situation returned).
2178        if matches!(
2179            status,
2180            RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
2181        ) {
2182            if let Some(key) = g.str_field("dedup_key") {
2183                set.insert(key.to_string());
2184            }
2185        }
2186    }
2187    Ok(set)
2188}
2189
2190fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
2191    let g = sub
2192        .grain(rec_hash)?
2193        .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
2194    Recommendation::from_fields(rec_hash, &g.fields)
2195}
2196
2197/// Is this line a definition rewrite — a statement that changes a saved
2198/// `qry:`/`tpl:` registry row rather than writing a grain?
2199///
2200/// A keyword test, not a parse: the engine deliberately contains a CAL
2201/// *writer*, never a parser (parsing is the substrate's job). Both spellings
2202/// are matched case-insensitively, and `DROP` is intentionally absent — the
2203/// loop may propose defining a query, never removing one.
2204pub(crate) fn is_definition_statement(line: &str) -> bool {
2205    let up = line.trim_start().to_ascii_uppercase();
2206    up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
2207}
2208
2209/// The refusal an advisory `Edit` earns. The engine has no executable edit
2210/// primitive; the change belongs in the host.
2211const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
2212     Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
2213     approve it to acknowledge it and let it expire.";
2214
2215/// The refusal an advisory `Data` finding earns — every `Data` shape except
2216/// `outcome_review`'s revert, which carries `revert_of`.
2217/// Shared by [`Engine::preflight_apply`] and the apply gate so a fused
2218/// approve-and-apply caller is refused BEFORE the approval lands, not after.
2219const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
2220     (evalset hash + run id + stats) — use apply_gated";
2221
2222const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
2223     guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
2224     acknowledge it and let it expire.";
2225
2226/// Whether [`Engine::apply`] can execute this proposal at all.
2227///
2228/// One source of truth, shared by [`Engine::preflight_apply`] and
2229/// [`Engine::apply`] so the two can never disagree.
2230///
2231/// Deliberately **not** consulted on approve. An advisory finding is a `Flag`:
2232/// approving it means "yes, this is real", which is the whole workflow for the
2233/// LLM path and the telemetry analyzers. What must not happen is a caller being
2234/// walked into an approval and *then* refused — which is exactly what the fused
2235/// approve-and-apply path in the bindings did, leaving the recommendation in
2236/// `approved`, whose only exits are `applied` and `expired`. Preflight asks
2237/// first, so that path now refuses before it commits anything.
2238pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
2239    match proposal {
2240        Proposal::Cal { .. } => Ok(()),
2241        Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
2242        // Two executable Data shapes: `outcome_review`'s `revert_of` (names
2243        // an earlier applied recommendation to roll back), and a gated
2244        // revision (code or adapter), which executes by writing its
2245        // promotion grain — the gating-run requirement itself is checked at
2246        // apply, not here.
2247        Proposal::Data { data } => {
2248            if requires_gating(action_kind)
2249                || data.get("revert_of").and_then(Value::as_str).is_some()
2250            {
2251                Ok(())
2252            } else {
2253                Err(Error::InvalidProposal(ADVISORY_DATA.into()))
2254            }
2255        }
2256    }
2257}
2258
2259#[cfg(test)]
2260mod definition_proposal_tests {
2261    use super::is_definition_statement;
2262
2263    #[test]
2264    fn definition_statements_are_recognized_in_both_spellings() {
2265        assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
2266        assert!(is_definition_statement("  define template foo AS { x }"));
2267        assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
2268        // Ordinary proposals are untouched.
2269        assert!(!is_definition_statement("ADD fact {}"));
2270        assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
2271        assert!(!is_definition_statement("FORGET abc"));
2272        // `DROP` is never proposable, so it is deliberately NOT a definition
2273        // statement here — a proposal containing one still fails validation
2274        // rather than being handed an inverse.
2275        assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
2276    }
2277}