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