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    pub fn preflight_apply<S: OmsSubstrate>(
1026        &self,
1027        sub: &S,
1028        rec_hash: &str,
1029        scopes: &ScopeSet,
1030        allow_destructive: bool,
1031    ) -> Result<()> {
1032        if !scopes.has(Scope::Apply) {
1033            return Err(Error::ScopeDenied("apply".into()));
1034        }
1035        let rec = load_rec(sub, rec_hash)?;
1036        if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1037            return Err(Error::DestructiveGated(
1038                "destructive apply requires admin scope + allow_destructive".into(),
1039            ));
1040        }
1041        ensure_executable(&rec.proposal)?;
1042        Ok(())
1043    }
1044
1045    /// Apply an approved recommendation. Requires `apply`; destructive payloads
1046    /// additionally require `admin` + `allow_destructive`. Records the applied
1047    /// info (inverse plan) for rollback.
1048    #[allow(clippy::too_many_arguments)]
1049    pub fn apply<S: OmsSubstrate>(
1050        &self,
1051        sub: &mut S,
1052        rec_hash: &str,
1053        actor: &str,
1054        observer: ObserverType,
1055        scopes: &ScopeSet,
1056        because: &str,
1057        allow_destructive: bool,
1058        now_ms: i64,
1059    ) -> Result<AppliedRecord> {
1060        self.apply_inner(
1061            sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1062        )
1063    }
1064
1065    /// Apply WITH the §7.4 evalset-run edge — the only path that can apply a
1066    /// `code_revision`. The evidence is validated against the
1067    /// recommendation's pin and recorded on the audit Observation.
1068    #[allow(clippy::too_many_arguments)]
1069    pub fn apply_gated<S: OmsSubstrate>(
1070        &self,
1071        sub: &mut S,
1072        rec_hash: &str,
1073        actor: &str,
1074        observer: ObserverType,
1075        scopes: &ScopeSet,
1076        because: &str,
1077        allow_destructive: bool,
1078        gating: &crate::recommendation::GatingEvidence,
1079        now_ms: i64,
1080    ) -> Result<AppliedRecord> {
1081        self.apply_inner(
1082            sub,
1083            rec_hash,
1084            actor,
1085            observer,
1086            scopes,
1087            because,
1088            allow_destructive,
1089            Some(gating),
1090            now_ms,
1091        )
1092    }
1093
1094    #[allow(clippy::too_many_arguments)]
1095    fn apply_inner<S: OmsSubstrate>(
1096        &self,
1097        sub: &mut S,
1098        rec_hash: &str,
1099        actor: &str,
1100        observer: ObserverType,
1101        scopes: &ScopeSet,
1102        because: &str,
1103        allow_destructive: bool,
1104        gating: Option<&crate::recommendation::GatingEvidence>,
1105        now_ms: i64,
1106    ) -> Result<AppliedRecord> {
1107        if !scopes.has(Scope::Apply) {
1108            return Err(Error::ScopeDenied("apply".into()));
1109        }
1110        let because = validate_because(because)?;
1111        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1112        let status = *p
1113            .status_index
1114            .get(rec_hash)
1115            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1116        if !status.can_transition_to(RecStatus::Applied, false) {
1117            return Err(Error::LifecycleViolation(format!(
1118                "{} -> applied (approve first)",
1119                status.as_str()
1120            )));
1121        }
1122        let rec = load_rec(sub, rec_hash)?;
1123        if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1124            return Err(Error::DestructiveGated(
1125                "destructive apply requires admin scope + allow_destructive".into(),
1126            ));
1127        }
1128        // §7.4: a code revision applies ONLY through the evalset-run edge —
1129        // the pin must match, the pinned evalset must still be LIVE (a
1130        // superseded evalset invalidates in-flight code recommendations:
1131        // they re-gate), and a failing gate admits nothing.
1132        if rec.action_kind == ActionKind::CodeRevision {
1133            let g = gating.ok_or_else(|| {
1134                Error::InvalidProposal(
1135                    "code revisions apply only with a recorded gating run \
1136                     (evalset hash + run id + stats) — use apply_gated"
1137                        .into(),
1138                )
1139            })?;
1140            let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1141            if g.evalset_hash != pin {
1142                return Err(Error::InvalidProposal(format!(
1143                    "gating ran evalset {} but the recommendation is pinned \
1144                     to {pin} (Rule E1)",
1145                    g.evalset_hash
1146                )));
1147            }
1148            match sub.grain(pin)? {
1149                Some(evalset) if evalset.is_live() => {}
1150                Some(_) => {
1151                    return Err(Error::InvalidProposal(
1152                        "the pinned evalset was superseded after gating — \
1153                         the recommendation must re-gate (Rule E1)"
1154                            .into(),
1155                    ))
1156                }
1157                None => {
1158                    return Err(Error::InvalidProposal(format!(
1159                        "pinned evalset {pin} not found in the substrate"
1160                    )))
1161                }
1162            }
1163            if g.failed > 0 {
1164                return Err(Error::InvalidProposal(format!(
1165                    "the gating run failed {}/{} cases — a failing gate \
1166                     cannot admit code",
1167                    g.failed,
1168                    g.passed + g.failed
1169                )));
1170            }
1171        }
1172
1173        // Execute the proposal.
1174        let mut created = Vec::new();
1175        // The inverse of a change that creates no grain — see
1176        // `AppliedRecord::inverse_cal`. Captured BEFORE execution, because
1177        // afterwards the previous definition is gone.
1178        let mut inverse_cal: Option<String> = None;
1179        match &rec.proposal {
1180            Proposal::Cal { cal } => {
1181                for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1182                    if !is_definition_statement(line) {
1183                        continue;
1184                    }
1185                    match sub.definition_inverse(line)? {
1186                        Some(inv) => inverse_cal = Some(inv),
1187                        None => {
1188                            return Err(Error::InvalidProposal(format!(
1189                                "this substrate cannot record a rollback inverse for {line:?}; \
1190                                 a definition rewrite that ROLLBACK could not undo is refused \
1191                                 rather than applied"
1192                            )))
1193                        }
1194                    }
1195                }
1196                let rows = sub.execute_cal(cal)?;
1197                for r in rows {
1198                    if let Some(h) = r.get("hash").and_then(Value::as_str) {
1199                        created.push(h.to_string());
1200                    }
1201                }
1202            }
1203            // The engine has no executable Edit primitive. Marking this
1204            // Applied used to be a lie (and rollback had no inverse).
1205            Proposal::Edit { .. } => {
1206                return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1207            }
1208            // A gated code revision EXECUTES by writing the promotion grain:
1209            // an immutable record that this tool target now resolves to the
1210            // proposed code (the payload rides the recommendation's Data —
1211            // typically a blob address from the §7.4 seam). Hosts read the
1212            // promotion to re-resolve; retracting it is the rollback
1213            // inverse, so the apply is rollbackable end-to-end.
1214            Proposal::Data { data } if rec.action_kind == ActionKind::CodeRevision => {
1215                let mut spec = crate::substrate::GrainSpec::new(
1216                    crate::model::grain_type::FACT,
1217                    LOOP_NS,
1218                )
1219                .with_field("subject", rec.target_ref.clone())
1220                .with_field("relation", "mg:code_promotion")
1221                .with_field(
1222                    "object",
1223                    serde_json::to_string(data)
1224                        .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1225                )
1226                .with_field("rec_hash", rec_hash.to_string());
1227                if let Some(g) = gating {
1228                    spec = spec
1229                        .with_field("gating_evalset", g.evalset_hash.clone())
1230                        .with_field("gating_run_id", g.run_id.clone());
1231                }
1232                created.push(sub.put_grain(&spec)?);
1233            }
1234            Proposal::Data { data } => {
1235                // OutcomeReview is the one executable Data shape: its
1236                // `revert_of` points at an earlier applied recommendation.
1237                // Reuse the ordinary rollback path so the created hashes are
1238                // really retracted and the original lifecycle/audit advances.
1239                let revert_of = data
1240                    .get("revert_of")
1241                    .and_then(Value::as_str)
1242                    .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1243                self.rollback(
1244                    sub,
1245                    revert_of,
1246                    actor,
1247                    observer,
1248                    scopes,
1249                    &because,
1250                    now_ms,
1251                )?;
1252                // rollback stored a newer lifecycle state; merge this apply
1253                // into that state rather than overwriting the rollback.
1254                p = LoopPersisted::from_value(sub.load_state()?)?;
1255            }
1256        }
1257
1258        let applied = AppliedRecord {
1259            applied_at_ms: now_ms,
1260            target_ref: rec.target_ref.clone(),
1261            rollbackable: rec.rollbackable,
1262            created_hashes: created,
1263            inverse_cal,
1264            metric: rec.metric.clone(),
1265        };
1266        let prev = p.audit_heads.get(rec_hash).cloned();
1267        let audit = AuditRecord {
1268            rec_hash: rec_hash.into(),
1269            from: Some(status),
1270            to: RecStatus::Applied,
1271            actor: actor.into(),
1272            observer_type: observer,
1273            because,
1274            previous_audit_hash: prev,
1275            gating: gating.cloned(),
1276            at_ms: now_ms,
1277        };
1278        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1279        p.audit_heads.insert(rec_hash.into(), audit_hash);
1280        p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1281        p.applied.insert(rec_hash.into(), applied.clone());
1282        sub.store_state(&p.to_value()?)?;
1283        Ok(applied)
1284    }
1285
1286    /// Roll back an applied recommendation by retracting the grains it created.
1287    /// Fails for non-rollbackable applies (e.g. FORGET).
1288    #[allow(clippy::too_many_arguments)]
1289    pub fn rollback<S: OmsSubstrate>(
1290        &self,
1291        sub: &mut S,
1292        rec_hash: &str,
1293        actor: &str,
1294        observer: ObserverType,
1295        scopes: &ScopeSet,
1296        because: &str,
1297        now_ms: i64,
1298    ) -> Result<()> {
1299        if !scopes.has(Scope::Apply) {
1300            return Err(Error::ScopeDenied("apply".into()));
1301        }
1302        let because = validate_because(because)?;
1303        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1304        let status = *p
1305            .status_index
1306            .get(rec_hash)
1307            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1308        if !status.can_transition_to(RecStatus::RolledBack, false) {
1309            return Err(Error::LifecycleViolation(format!(
1310                "{} -> rolled_back",
1311                status.as_str()
1312            )));
1313        }
1314        let applied = p
1315            .applied
1316            .get(rec_hash)
1317            .cloned()
1318            .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1319        if !applied.rollbackable {
1320            return Err(Error::LifecycleViolation(
1321                "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1322            ));
1323        }
1324        for h in &applied.created_hashes {
1325            sub.retract(h, &format!("rollback of {rec_hash}"))?;
1326        }
1327        // A definition rewrite creates no grain, so retracting `created_hashes`
1328        // undoes nothing. Restoring it means re-running the statement captured
1329        // at apply time — the previous definition, or a DROP when there was
1330        // none. Runs BEFORE the audit is written, so a failed restore leaves
1331        // the recommendation `applied` (still true) rather than recording a
1332        // rollback that did not happen.
1333        if let Some(inverse) = &applied.inverse_cal {
1334            sub.execute_cal(inverse)?;
1335        }
1336        let prev = p.audit_heads.get(rec_hash).cloned();
1337        let audit = AuditRecord {
1338            rec_hash: rec_hash.into(),
1339            from: Some(status),
1340            to: RecStatus::RolledBack,
1341            actor: actor.into(),
1342            observer_type: observer,
1343            because,
1344            previous_audit_hash: prev,
1345            gating: None,
1346            at_ms: now_ms,
1347        };
1348        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1349        p.audit_heads.insert(rec_hash.into(), audit_hash);
1350        p.status_index
1351            .insert(rec_hash.into(), RecStatus::RolledBack);
1352        sub.store_state(&p.to_value()?)?;
1353        Ok(())
1354    }
1355
1356    /// List stored recommendations, optionally filtered by status. Status comes
1357    /// from the rebuildable index, not the immutable grain body. Ordered for
1358    /// review triage — highest severity first, then oldest first — and stable
1359    /// across runs for identical input.
1360    pub fn recommendations<S: OmsSubstrate>(
1361        &self,
1362        sub: &S,
1363        status_filter: Option<RecStatus>,
1364    ) -> Result<Vec<Recommendation>> {
1365        let p = LoopPersisted::from_value(sub.load_state()?)?;
1366        let grains = sub.grains_of_type(
1367            crate::model::grain_type::RECOMMENDATION,
1368            Some(LOOP_NS),
1369            ReadOpts {
1370                live_only: false,
1371                since_ms: None,
1372            },
1373        )?;
1374        let mut out = Vec::new();
1375        for g in grains {
1376            let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1377            rec.status = p
1378                .status_index
1379                .get(&g.hash)
1380                .copied()
1381                .unwrap_or(RecStatus::Pending);
1382            if let Some(f) = status_filter {
1383                if rec.status != f {
1384                    continue;
1385                }
1386            }
1387            out.push(rec);
1388        }
1389        // Review-queue order: worst first, then oldest first. Hash is only the
1390        // final tiebreak — sorting by it alone is deterministic per run but
1391        // meaningless across runs, because a grain's hash covers its timestamp,
1392        // so an identical queue comes back in a different order every time.
1393        // `dedup_key` is the last tiebreak that actually decides anything: it is
1394        // content-derived and stable across runs, whereas findings proposed in
1395        // the same sweep routinely share a `created_at_ms`. Hash trails it only
1396        // to make the ordering total.
1397        out.sort_by(|a, b| {
1398            b.severity
1399                .cmp(&a.severity)
1400                .then(a.created_at_ms.cmp(&b.created_at_ms))
1401                .then(a.dedup_key.cmp(&b.dedup_key))
1402                .then(a.hash.cmp(&b.hash))
1403        });
1404        Ok(out)
1405    }
1406
1407    /// Per-analyzer effective settings for the Setup view: the manifest facts
1408    /// merged with the file-config (override or manifest default). Read-only.
1409    pub fn analyzer_settings<S: OmsSubstrate>(
1410        &self,
1411        sub: &S,
1412    ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1413        let p = LoopPersisted::from_value(sub.load_state()?)?;
1414        Ok(self
1415            .analyzers
1416            .iter()
1417            .map(|a| {
1418                let m = a.manifest();
1419                let cfg = p.config.get(&m.id);
1420                crate::config::AnalyzerSetting {
1421                    id: m.id.clone(),
1422                    title: m.title.clone(),
1423                    description: m.description.clone(),
1424                    tier: format!("{:?}", m.tier),
1425                    trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1426                    default_on: m.default_on,
1427                    enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1428                    severity_floor: cfg
1429                        .and_then(|c| c.severity_floor)
1430                        .map(|s| s.as_str().to_string()),
1431                }
1432            })
1433            .collect())
1434    }
1435
1436    /// Update one analyzer's file-config (enable/disable, severity floor, param
1437    /// overrides, namespace scoping). Requires `Admin`. Params are validated
1438    /// against the analyzer's manifest first (unknown keys rejected, fail-closed),
1439    /// and the analyzer must exist. Returns the merged config as stored. This is
1440    /// the only write into `persisted.config` — the config layer, never a grain.
1441    pub fn set_analyzer_config<S: OmsSubstrate>(
1442        &self,
1443        sub: &mut S,
1444        analyzer_id: &str,
1445        update: crate::config::AnalyzerConfigUpdate,
1446        scopes: &ScopeSet,
1447    ) -> Result<crate::config::AnalyzerConfig> {
1448        if !scopes.has(Scope::Admin) {
1449            return Err(Error::ScopeDenied("admin".into()));
1450        }
1451        let manifest = self
1452            .analyzers
1453            .iter()
1454            .map(|a| a.manifest())
1455            .find(|m| m.id == analyzer_id)
1456            .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1457        // Validate params against the manifest BEFORE touching state.
1458        if let Some(params) = &update.params {
1459            manifest.resolve_params(params)?;
1460        }
1461        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1462        let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1463        if let Some(enabled) = update.enabled {
1464            cfg.enabled = Some(enabled);
1465        }
1466        if update.clear_floor {
1467            cfg.severity_floor = None;
1468        } else if let Some(floor) = update.severity_floor {
1469            cfg.severity_floor = Some(floor);
1470        }
1471        if let Some(params) = update.params {
1472            cfg.params = params;
1473        }
1474        if let Some(ns) = update.namespaces {
1475            cfg.namespaces = ns;
1476        }
1477        let stored = cfg.clone();
1478        sub.store_state(&p.to_value()?)?;
1479        Ok(stored)
1480    }
1481
1482    /// The measured outcome time series (the Verify gate's history) across all
1483    /// recommendations, ordered by when each checkpoint was measured.
1484    pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1485        let p = LoopPersisted::from_value(sub.load_state()?)?;
1486        let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1487        // `metric` and `rec_hash` break the tie: checkpoints measured in the
1488        // same sweep share a `measured_at_ms`, and without a tiebreak the order
1489        // falls through to the map's rec_hash ordering, which shifts every run.
1490        out.sort_by(|a, b| {
1491            a.measured_at_ms
1492                .cmp(&b.measured_at_ms)
1493                .then(a.horizon_ms.cmp(&b.horizon_ms))
1494                .then(a.metric.cmp(&b.metric))
1495                .then(a.rec_hash.cmp(&b.rec_hash))
1496        });
1497        Ok(out)
1498    }
1499
1500    /// A health snapshot — when the loop last ran, how much is un-analyzed
1501    /// since, and the queue counts. Lets a host surface "the loop may be stale"
1502    /// so a forgotten SessionEnd hook / cron doesn't silently kill it.
1503    pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1504        let p = LoopPersisted::from_value(sub.load_state()?)?;
1505        let (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
1506        let recs = self.recommendations(sub, None)?;
1507        let mut pending = 0;
1508        let mut applied = 0;
1509        for r in &recs {
1510            match r.status {
1511                RecStatus::Pending => pending += 1,
1512                RecStatus::Applied => applied += 1,
1513                _ => {}
1514            }
1515        }
1516        // Stale if it has never run, or it's been a while / a lot has piled up.
1517        let stale = match p.state.last_run_ms {
1518            None => true,
1519            Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1520        };
1521        Ok(Health {
1522            last_run_ms: p.state.last_run_ms,
1523            grains_since_run,
1524            error_events_since_run,
1525            pending,
1526            applied,
1527            total: recs.len() as u64,
1528            stale,
1529        })
1530    }
1531
1532    /// Approval-rate metric for `origin = llm` recommendations (reflection
1533    /// design §6b) — the live field-quality signal that accrues off the audit
1534    /// chain: what fraction of the model's *surfaced* proposals a reviewer
1535    /// accepts. Complements the offline Effective-Reliability eval.
1536    pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1537        let recs = self.recommendations(sub, None)?;
1538        let mut m = LlmMetrics::default();
1539        for r in &recs {
1540            if !matches!(r.origin, Origin::Llm { .. }) {
1541                continue;
1542            }
1543            m.proposed += 1;
1544            match r.status {
1545                RecStatus::Pending => m.pending += 1,
1546                RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1547                RecStatus::Rejected => m.rejected += 1,
1548                RecStatus::Expired => {}
1549            }
1550        }
1551        let decided = m.approved + m.rejected;
1552        m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1553        Ok(m)
1554    }
1555}
1556
1557/// A health snapshot for the backend's self-improvement loop.
1558#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1559pub struct Health {
1560    #[serde(skip_serializing_if = "Option::is_none")]
1561    pub last_run_ms: Option<i64>,
1562    pub grains_since_run: u64,
1563    pub error_events_since_run: u64,
1564    pub pending: u64,
1565    pub applied: u64,
1566    pub total: u64,
1567    /// True when the loop looks stalled (never run, or ≥7d / ≥100 new grains
1568    /// since the last run) — a nudge that a trigger may be unwired.
1569    pub stale: bool,
1570}
1571
1572/// Approval-rate metric for `origin = llm` recommendations (reflection §6b).
1573#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1574pub struct LlmMetrics {
1575    /// Total llm-origin recommendations ever stored (those that survived the
1576    /// verifier and reached the queue).
1577    pub proposed: u64,
1578    pub pending: u64,
1579    /// Approved + Applied + RolledBack (a reviewer said yes at least once).
1580    pub approved: u64,
1581    pub rejected: u64,
1582    /// approved / (approved + rejected); `None` until at least one is decided.
1583    #[serde(skip_serializing_if = "Option::is_none")]
1584    pub approval_rate: Option<f64>,
1585}
1586
1587/// Re-measure applied recommendations at each **checkpoint** past due, via the
1588/// engine's typed reads — no CAL-scalar round-trip. A recommendation
1589/// accumulates one `OutcomeResult` per horizon (measured once each), forming a
1590/// time series, so a late regression (held at 1d, regressed at 30d) is caught.
1591/// Only *regressed* checkpoints feed the outcome analyzer (→ a revert).
1592/// Unknown metric kinds are skipped, never faked.
1593fn measure_outcomes<S: OmsSubstrate>(
1594    sub: &S,
1595    p: &mut LoopPersisted,
1596    now_ms: i64,
1597) -> Result<Vec<OutcomeInput>> {
1598    // Collect all due (recommendation, horizon) checkpoints first.
1599    let mut due: Vec<(String, crate::config::AppliedRecord, i64)> = Vec::new();
1600    for (h, a) in &p.applied {
1601        if p.status_index.get(h) != Some(&RecStatus::Applied) {
1602            continue;
1603        }
1604        let Some(metric) = &a.metric else { continue };
1605        let done = p.measured.get(h).cloned().unwrap_or_default();
1606        for horizon in metric.horizons() {
1607            if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
1608                due.push((h.clone(), a.clone(), horizon));
1609            }
1610        }
1611    }
1612
1613    let mut out = Vec::new();
1614    for (rec_hash, applied, horizon) in due {
1615        let metric = applied.metric.as_ref().unwrap();
1616        let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1617            continue; // metric kind not yet re-measurable
1618        };
1619        let regressed = crate::recommendation::is_regression(
1620            metric.baseline,
1621            current,
1622            metric.higher_is_better,
1623        );
1624        p.outcomes.entry(rec_hash.clone()).or_default().push(
1625            crate::recommendation::OutcomeResult {
1626                rec_hash: rec_hash.clone(),
1627                metric: metric.metric.clone(),
1628                baseline: metric.baseline,
1629                current,
1630                verdict: if regressed { "regressed" } else { "held" }.into(),
1631                horizon_ms: horizon,
1632                measured_at_ms: now_ms,
1633            },
1634        );
1635        p.measured.entry(rec_hash.clone()).or_default().push(horizon);
1636        if regressed {
1637            out.push(OutcomeInput {
1638                rec_hash,
1639                target_ref: applied.target_ref.clone(),
1640                metric: metric.metric.clone(),
1641                baseline: metric.baseline,
1642                current,
1643                unit: metric.unit.clone(),
1644                higher_is_better: metric.higher_is_better,
1645            });
1646        }
1647    }
1648    Ok(out)
1649}
1650
1651/// Typed re-measurement for the fixed set of metric kinds the engine knows.
1652pub(crate) fn measure_metric<S: SubstrateRead>(
1653    sub: &S,
1654    metric: &crate::recommendation::MetricSnapshot,
1655    since_ms: i64,
1656) -> Result<Option<f64>> {
1657    match metric.metric.as_str() {
1658        // How many times did this tool fail again *with the same signature*
1659        // after the lesson was applied? Scoped to the signature (metric.relation)
1660        // so an unrelated later failure of the same tool is not read as a
1661        // regression of this specific lesson.
1662        "tool_error_recurrence" => {
1663            let Some(tool) = &metric.subject else { return Ok(None) };
1664            let tools = sub.grains_of_type(
1665                crate::model::grain_type::TOOL,
1666                None,
1667                ReadOpts { live_only: true, since_ms: Some(since_ms) },
1668            )?;
1669            let n = tools
1670                .iter()
1671                .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
1672                .filter(|t| {
1673                    // No stored signature (legacy metric) → fall back to the
1674                    // whole-tool count so old recommendations still measure.
1675                    metric.relation.as_deref().is_none_or(|sig| {
1676                        crate::analyzers::tool_failure::normalize_signature(
1677                            t.tool_content().unwrap_or(""),
1678                        ) == sig
1679                    })
1680                })
1681                .count();
1682            Ok(Some(n as f64))
1683        }
1684        // After a resolve-to-latest, does the subject again hold more than one
1685        // live value under the functional relation? Live-state read (no since
1686        // filter): the excess beyond one distinct object is the regression.
1687        "contradiction_recurrence" => {
1688            let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
1689                return Ok(None);
1690            };
1691            let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
1692            let distinct: BTreeSet<String> = facts
1693                .iter()
1694                .filter(|f| {
1695                    f.fact_relation()
1696                        .is_some_and(|r| normalize_ident(r) == *relation)
1697                })
1698                .filter_map(|f| f.fact_object().map(normalize_ident))
1699                .collect();
1700            Ok(Some(distinct.len().saturating_sub(1) as f64))
1701        }
1702        // `evalset:<hash>:<field>` — external correctness, measured the only
1703        // way that keeps the honesty rule ("Areev Loop improves the agent's
1704        // memory, not its outputs") intact: an evalset run is an INTERNAL,
1705        // BOUNDED, ATTRIBUTABLE measurement. The engine never runs the evalset;
1706        // it only reads summaries a host journaled with `areev eval run`.
1707        //
1708        // `since_ms` is the apply time, and it is load-bearing rather than an
1709        // optimization: a run journaled BEFORE the apply cannot be evidence of
1710        // what applying did. Comparing the baseline run against itself would
1711        // report "held" forever — a fabricated receipt, which is worse than no
1712        // receipt at all. No run since the apply → `None` → not yet
1713        // measurable, and the checkpoint stays due.
1714        m if m.starts_with("evalset:") => {
1715            let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
1716                return Ok(None);
1717            };
1718            let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
1719                return Ok(None);
1720            };
1721            // `failed`/`passed`/`total` are promoted so a metric can be written
1722            // against any evalset without the host having to add fields;
1723            // anything else is read from the summary the host did write.
1724            Ok(match field {
1725                "failed" => Some(run.failed as f64),
1726                "passed" => Some(run.passed as f64),
1727                "total" => Some(run.total() as f64),
1728                "error_rate" => match run.total() {
1729                    0 => None, // no cases ran: undefined, not zero
1730                    t => Some(run.failed as f64 / t as f64),
1731                },
1732                other => run.field(other),
1733            })
1734        }
1735        _ => Ok(None),
1736    }
1737}
1738
1739/// Live facts for one normalized (namespace?, subject) — the shared scope of
1740/// the fact-shaped recurrence metrics. `namespace: None` spans all namespaces.
1741fn scoped_live_facts<S: SubstrateRead>(
1742    sub: &S,
1743    namespace: Option<&str>,
1744    subject: &str,
1745) -> Result<Vec<GrainRecord>> {
1746    let facts = sub.grains_of_type(
1747        crate::model::grain_type::FACT,
1748        None,
1749        ReadOpts { live_only: true, since_ms: None },
1750    )?;
1751    Ok(facts
1752        .into_iter()
1753        .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
1754        .filter(|f| {
1755            f.fact_subject()
1756                .is_some_and(|s| normalize_ident(s) == subject)
1757        })
1758        .collect())
1759}
1760
1761// --- free helpers ---
1762
1763/// The §6.3 exact-equality check: a SUPERSEDE is *value-identical* when every
1764/// replacement field equals the superseded grain's value — strings after
1765/// case-fold/trim (upstream NFC is an OMS invariant), `namespace` against the
1766/// grain's own namespace, everything else exactly. This is what makes an
1767/// auto-applied consolidation provably information-preserving; a
1768/// near-duplicate (an observation body off by one token) fails it and stays
1769/// pending for human review. Fails closed: an unrecognized line shape, an
1770/// empty replacement, a missing grain, or a field the original never had all
1771/// disqualify.
1772fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
1773    let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
1774        return false;
1775    };
1776    if fields.is_empty() {
1777        return false;
1778    }
1779    let Ok(Some(grain)) = sub.grain(&target) else {
1780        return false;
1781    };
1782    // An expiry lives OUTSIDE `fields`, so the replacement (built from a fixed
1783    // field set) can never carry it — consolidating away a grain that has a
1784    // valid_to would silently drop the expiry, invisibly to the field-by-field
1785    // check below. Fail closed. (A dup that additionally carries an extra
1786    // *content* field the replacement omits is a real but subtler info-loss;
1787    // catching it soundly needs a full canonical-vs-extra comparison rather
1788    // than this replacement-scoped check, since a real fact's fields also carry
1789    // OMS metadata like `confidence` that consolidation legitimately keeps —
1790    // left as a follow-up so this narrow fix can't block valid consolidations.)
1791    if grain.valid_to_ms.is_some() {
1792        return false;
1793    }
1794    // Forward check: every replacement field equals the grain's value.
1795    fields.iter().all(|(k, v)| {
1796        if k == "namespace" {
1797            return v
1798                .as_str()
1799                .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
1800        }
1801        match (v, grain.fields.get(k)) {
1802            (Value::String(a), Some(Value::String(b))) => {
1803                normalize_ident(a) == normalize_ident(b)
1804            }
1805            (a, Some(b)) => a == b,
1806            (_, None) => false,
1807        }
1808    })
1809}
1810
1811/// The confidence floor (§5.4): a verified draft below this is dropped. The
1812/// verifier's calibrated confidence is the gate, not the proposer's self-report.
1813const MIN_LLM_CONFIDENCE: f64 = 0.75;
1814
1815/// The fixed DISCOVER instruction (§5.1). The scoring rule makes "nothing to
1816/// report" a first-class, zero-penalty answer — the structural antidote to
1817/// over-generation. Kept in its own request field so it never interleaves with
1818/// (attacker-influenced) evidence text.
1819const DISCOVER_INSTRUCTIONS: &str = "You review an agent's memory for quality. \
1820Given deterministic findings and the evidence they cite, propose ADDITIONAL \
1821findings the deterministic checks would miss (e.g. a semantic contradiction, a \
1822stale assumption, a duplicated meaning). SCORING: propose a finding ONLY if you \
1823are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
1824useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
1825earns 0. When in doubt, propose nothing — an empty list is the correct answer \
1826when there is nothing worth flagging. The 'approved' and 'rejected' lists, when \
1827present, show findings this reviewer recently accepted or rejected — prefer the \
1828kind they accept and avoid the kind they reject. Every proposal MUST cite one \
1829or more evidence hashes from the bundle, target a memory entity, and include \
1830your confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
1831\"target\":\"entity:<ns>/<subject>\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
1832\"confidence\":0.0}]}. Propose nothing you cannot ground in the evidence.";
1833
1834/// The fixed GROUND instruction (§5.2): verify the finding's factual PREMISES
1835/// are real (anti-fabrication), while allowing an inference. A self-improvement
1836/// finding may reason BEYOND the evidence (e.g. 'HQ=SF and country=Germany are
1837/// inconsistent'); grounding checks the premises (HQ=SF, country=Germany) are in
1838/// the evidence — the soundness of the inference is VERIFY's job, not this one.
1839const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
1840fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
1841the facts it relies on are actually present in the cited evidence, NOT that its \
1842conclusion is stated verbatim. Decompose the finding into the factual claims it \
1843depends on. Mark supported=true when those facts are present in the evidence \
1844(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
1845on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
1846different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
1847
1848/// The fixed VERIFY instruction (§5.3): adversarial, abstention-biased.
1849const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
1850each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
1851never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
1852SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
1853'possible' findings with no concrete defect, and reject any claimed \
1854inconsistency or contradiction that is not backed by at least two actually \
1855conflicting facts in the cited evidence. (2) Context — does the finding \
1856correctly read its cited evidence, or misinterpret what the grains say? Keep a \
1857finding when it names a genuine, specific problem grounded in its evidence and \
1858materially useful to a human reviewer; otherwise reject it, and default to \
1859keep=false when uncertain. Do NOT reject a finding for being 'already known', \
1860redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
1861grounded cross-fact inconsistency with two conflicting facts is exactly what to \
1862KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
1863{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
1864
1865/// The fixed ENRICH instruction.
1866const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
1867guidance note to help a human reviewer decide. Do not restate the finding. Return \
1868JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
1869
1870/// Add a grain to the evidence bundle — deduplicated, bounded at 64, text
1871/// capped. Shared by the deterministic-citation and recent-grain seeding.
1872fn push_evidence(
1873    evidence: &mut Vec<crate::llm::EvidenceItem>,
1874    bundle: &mut BTreeSet<String>,
1875    g: &GrainRecord,
1876) {
1877    if evidence.len() < 64 && bundle.insert(g.hash.clone()) {
1878        evidence.push(crate::llm::EvidenceItem {
1879            hash: g.hash.clone(),
1880            grain_type: g.grain_type.clone(),
1881            text: crate::llm::cap(&grain_brief(g), 400),
1882        });
1883    }
1884}
1885
1886/// A short human-readable projection of a grain for the evidence bundle.
1887fn grain_brief(g: &GrainRecord) -> String {
1888    if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
1889        return format!("{s} {r} {o}");
1890    }
1891    for key in ["content", "body", "text", "summary"] {
1892        if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
1893            if !v.is_empty() {
1894                return v.to_string();
1895            }
1896        }
1897    }
1898    String::new()
1899}
1900
1901/// Stamp a validated DISCOVER draft as an `origin = llm` recommendation. LLM
1902/// drafts are always advisory `Flag`s carrying `Proposal::Data` (never an
1903/// executable CAL mutation), lower-confidence, and — via `Origin::Llm` and a
1904/// no-manifest analyzer id — structurally ineligible for auto-apply.
1905fn stamp_llm(
1906    model: &str,
1907    d: &crate::llm::LlmDraft,
1908    target_ref: String,
1909    cited: Vec<String>,
1910    confidence: f64,
1911    now_ms: i64,
1912) -> Recommendation {
1913    let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
1914    let mut args = serde_json::Map::new();
1915    args.insert("text".into(), Value::from(summary_text));
1916    let guidance = if d.guidance.trim().is_empty() {
1917        None
1918    } else {
1919        Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
1920    };
1921    let action = ActionKind::Flag;
1922    let mut data = serde_json::Map::new();
1923    data.insert("source".into(), Value::from("llm"));
1924    Recommendation {
1925        hash: String::new(),
1926        analyzer: "loop.llm/1".to_string(),
1927        params_snapshot: serde_json::Map::new(),
1928        origin: Origin::Llm { model: model.to_string() },
1929        target_ref: target_ref.clone(),
1930        action_kind: action,
1931        dedup_key: dedup_key("llm", &target_ref, action),
1932        summary: Summary::new("llm.discover", args),
1933        severity: Severity::Low,
1934        proposal: Proposal::Data { data },
1935        destructive: false,
1936        rollbackable: false,
1937        evidence: cited,
1938        evidence_query: None,
1939        metric: None,
1940        // The verifier's calibrated confidence — not a hardcoded default.
1941        confidence: confidence.clamp(0.0, 1.0),
1942        importance: 0.3,
1943        created_at_ms: now_ms,
1944        guidance,
1945        evalset_hash: None,
1946        status: RecStatus::Pending,
1947    }
1948}
1949
1950fn stamp(
1951    m: &AnalyzerManifest,
1952    params: &crate::manifest::Params,
1953    d: crate::recommendation::RecDraft,
1954    now_ms: i64,
1955) -> Result<Recommendation> {
1956    let target = TargetRef::parse(&d.target_ref)?;
1957    // Rule E1 at the door (§7.4): an unpinned code revision, a pinned
1958    // non-code target, or a code action on a non-tool target never becomes
1959    // a recommendation at all.
1960    crate::recommendation::validate_code_rules(
1961        d.action_kind,
1962        target.target_class(),
1963        d.evalset_hash.as_deref(),
1964    )?;
1965    let dedup = dedup_key(m.family(), &d.target_ref, d.action_kind);
1966    let destructive = match &d.proposal {
1967        Proposal::Cal { cal } => cal::contains_destructive(cal),
1968        _ => false,
1969    };
1970    let rollbackable = match &d.proposal {
1971        Proposal::Cal { .. } => !destructive,
1972        Proposal::Edit { .. } => false,
1973        // A code revision applies by WRITING the promotion grain; retracting
1974        // it is the exact inverse — rollbackable by construction.
1975        Proposal::Data { .. } => d.action_kind == ActionKind::CodeRevision,
1976    };
1977    let mut evidence = d.evidence;
1978    evidence.truncate(MAX_EVIDENCE);
1979    // Provenance follows the analyzer's trust class: a subprocess
1980    // (`--analyzer-cmd`) finding is stamped `Command` — structurally
1981    // auto-apply-ineligible and badged [external] on the recall surface —
1982    // not `Builtin`.
1983    let origin = match m.trust_class {
1984        crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
1985        _ => Origin::Builtin,
1986    };
1987    Ok(Recommendation {
1988        hash: String::new(),
1989        analyzer: m.id.clone(),
1990        params_snapshot: params.snapshot(),
1991        origin,
1992        target_ref: target.as_string(),
1993        action_kind: d.action_kind,
1994        dedup_key: dedup,
1995        summary: d.summary,
1996        severity: d.severity,
1997        proposal: d.proposal,
1998        destructive,
1999        rollbackable,
2000        evidence,
2001        evidence_query: d.evidence_query,
2002        metric: d.metric,
2003        confidence: d.confidence,
2004        importance: d.importance,
2005        created_at_ms: now_ms,
2006        guidance: None,
2007        evalset_hash: d.evalset_hash,
2008        status: RecStatus::Pending,
2009    })
2010}
2011
2012fn validate_because(because: &str) -> Result<String> {
2013    let trimmed = because.trim();
2014    if trimmed.is_empty() {
2015        return Err(Error::InvalidProposal(
2016            "a BECAUSE reason is required".into(),
2017        ));
2018    }
2019    if trimmed.chars().count() > MAX_BECAUSE {
2020        return Err(Error::InvalidProposal(format!(
2021            "BECAUSE exceeds {MAX_BECAUSE} chars"
2022        )));
2023    }
2024    Ok(trimmed.to_string())
2025}
2026
2027fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
2028    for req in &m.requires {
2029        match req {
2030            Capability::Forks if !caps.forks => return Some("forks"),
2031            Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
2032            Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
2033            _ => {}
2034        }
2035    }
2036    None
2037}
2038
2039fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
2040    p.config.get(analyzer_id).and_then(|c| c.severity_floor)
2041}
2042
2043fn gate(
2044    opts: &RunOptions,
2045    p: &LoopPersisted,
2046    new_grains: u64,
2047    new_errors: u64,
2048    now_ms: i64,
2049) -> Option<SkipReason> {
2050    let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
2051    if !any {
2052        return None;
2053    }
2054    let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
2055    let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
2056    let stale_ok = opts
2057        .if_stale_ms
2058        .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
2059    if min_new_ok || min_err_ok || stale_ok {
2060        return None;
2061    }
2062    // Choose the honest reason: staleness-only gate → not_stale, else min_new.
2063    if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
2064        Some(SkipReason::NotStale)
2065    } else {
2066        Some(SkipReason::MinNewNotMet)
2067    }
2068}
2069
2070fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
2071    let opts = ReadOpts {
2072        live_only: false,
2073        since_ms: watermark.map(|w| w + 1),
2074    };
2075    let mut new_grains = 0u64;
2076    let mut new_errors = 0u64;
2077    for t in [
2078        crate::model::grain_type::FACT,
2079        crate::model::grain_type::EVENT,
2080        crate::model::grain_type::TOOL,
2081        crate::model::grain_type::OBSERVATION,
2082    ] {
2083        let g = sub.grains_of_type(t, None, opts)?;
2084        new_grains += g.len() as u64;
2085        // The error gate (--min-new-errors) watches captured tool failures.
2086        if t == crate::model::grain_type::TOOL {
2087            new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
2088        }
2089    }
2090    Ok((new_grains, new_errors))
2091}
2092
2093fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
2094    let grains = sub.grains_of_type(
2095        crate::model::grain_type::RECOMMENDATION,
2096        Some(LOOP_NS),
2097        ReadOpts {
2098            live_only: false,
2099            since_ms: None,
2100        },
2101    )?;
2102    let mut set = BTreeSet::new();
2103    for g in grains {
2104        let status = p
2105            .status_index
2106            .get(&g.hash)
2107            .copied()
2108            .unwrap_or(RecStatus::Pending);
2109        // Pending/approved (still open) and applied (already handled)
2110        // recommendations suppress re-proposal of the same finding. Rejected
2111        // is handled by cooldowns; rolled_back/expired may legitimately
2112        // re-propose (the situation returned).
2113        if matches!(
2114            status,
2115            RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
2116        ) {
2117            if let Some(key) = g.str_field("dedup_key") {
2118                set.insert(key.to_string());
2119            }
2120        }
2121    }
2122    Ok(set)
2123}
2124
2125fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
2126    let g = sub
2127        .grain(rec_hash)?
2128        .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
2129    Recommendation::from_fields(rec_hash, &g.fields)
2130}
2131
2132/// Is this line a definition rewrite — a statement that changes a saved
2133/// `qry:`/`tpl:` registry row rather than writing a grain?
2134///
2135/// A keyword test, not a parse: the engine deliberately contains a CAL
2136/// *writer*, never a parser (parsing is the substrate's job). Both spellings
2137/// are matched case-insensitively, and `DROP` is intentionally absent — the
2138/// loop may propose defining a query, never removing one.
2139pub(crate) fn is_definition_statement(line: &str) -> bool {
2140    let up = line.trim_start().to_ascii_uppercase();
2141    up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
2142}
2143
2144/// The refusal an advisory `Edit` earns. The engine has no executable edit
2145/// primitive; the change belongs in the host.
2146const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
2147     Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
2148     approve it to acknowledge it and let it expire.";
2149
2150/// The refusal an advisory `Data` finding earns — every `Data` shape except
2151/// `outcome_review`'s revert, which carries `revert_of`.
2152const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
2153     guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
2154     acknowledge it and let it expire.";
2155
2156/// Whether [`Engine::apply`] can execute this proposal at all.
2157///
2158/// One source of truth, shared by [`Engine::preflight_apply`] and
2159/// [`Engine::apply`] so the two can never disagree.
2160///
2161/// Deliberately **not** consulted on approve. An advisory finding is a `Flag`:
2162/// approving it means "yes, this is real", which is the whole workflow for the
2163/// LLM path and the telemetry analyzers. What must not happen is a caller being
2164/// walked into an approval and *then* refused — which is exactly what the fused
2165/// approve-and-apply path in the bindings did, leaving the recommendation in
2166/// `approved`, whose only exits are `applied` and `expired`. Preflight asks
2167/// first, so that path now refuses before it commits anything.
2168pub(crate) fn ensure_executable(proposal: &Proposal) -> Result<()> {
2169    match proposal {
2170        Proposal::Cal { .. } => Ok(()),
2171        Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
2172        // `outcome_review` is the one executable Data shape: `revert_of` names
2173        // an earlier applied recommendation to roll back.
2174        Proposal::Data { data } => {
2175            if data.get("revert_of").and_then(Value::as_str).is_some() {
2176                Ok(())
2177            } else {
2178                Err(Error::InvalidProposal(ADVISORY_DATA.into()))
2179            }
2180        }
2181    }
2182}
2183
2184#[cfg(test)]
2185mod definition_proposal_tests {
2186    use super::is_definition_statement;
2187
2188    #[test]
2189    fn definition_statements_are_recognized_in_both_spellings() {
2190        assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
2191        assert!(is_definition_statement("  define template foo AS { x }"));
2192        assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
2193        // Ordinary proposals are untouched.
2194        assert!(!is_definition_statement("ADD fact {}"));
2195        assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
2196        assert!(!is_definition_statement("FORGET abc"));
2197        // `DROP` is never proposable, so it is deliberately NOT a definition
2198        // statement here — a proposal containing one still fails validation
2199        // rather than being handed an inverse.
2200        assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
2201    }
2202}