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    /// Where the LLM's contribution went, stage by stage. `None` when no
131    /// backend is attached.
132    #[serde(default, skip_serializing_if = "Option::is_none")]
133    pub llm_funnel: Option<LlmFunnel>,
134}
135
136/// The DISCOVER pipeline's attrition, counted.
137///
138/// "The model contributed nothing" has at least five distinct causes, and
139/// they call for opposite responses: an empty bundle is a capture problem, a
140/// model that abstained may need better evidence or a better prompt, drafts
141/// dying at GROUND suggest fabrication, drafts dying at VERIFY suggest they
142/// were vague, and drafts dying at the floor were merely unconfident. Without
143/// this they are indistinguishable from the outside — every one of them
144/// renders as an empty ledger and reads like a clean null. That ambiguity
145/// cost a six-cell measurement run before it was noticed.
146#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
147pub struct LlmFunnel {
148    /// Grains offered to the model as evidence. Zero means nothing to reflect on.
149    pub evidence: u64,
150    /// Drafts the model returned.
151    pub proposed: u64,
152    /// Survived the cite-check and target-class filter.
153    pub cited: u64,
154    /// Dropped for citing a hash the bundle does not hold. Its own counter
155    /// because the fix is specific: models copy long hex badly, and that is
156    /// a different problem from one aiming at a surface it may not touch.
157    pub dropped_uncited: u64,
158    /// Dropped for targeting a class the vocabulary may not reach — a prompt,
159    /// host config, or its own grader.
160    pub dropped_target: u64,
161    /// Survived GROUND — their premises were found in the cited evidence.
162    pub grounded: u64,
163    /// Verdicts GROUND actually returned. `grounded = 0` with verdicts > 0 is
164    /// the gate refusing every draft; `grounded = 0` with verdicts = 0 is a
165    /// grader that answered with nothing usable, which is a backend problem
166    /// wearing a gate's clothes.
167    pub ground_verdicts: u64,
168    /// The GROUND call itself failed — no response at all. The engine
169    /// fail-softs here by design, so without this the run looks like a model
170    /// that had nothing to say.
171    pub ground_call_failed: bool,
172    /// Survived VERIFY's adversarial pass.
173    pub kept: u64,
174    /// Cleared the confidence floor and reached the queue.
175    pub stored: u64,
176}
177
178impl RunResult {
179    fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
180        RunResult {
181            outcome: RunOutcome::Skipped,
182            skip_reason: Some(reason),
183            new_grains,
184            new_error_events,
185            proposed: 0,
186            deduped: 0,
187            stored: 0,
188            auto_applied: 0,
189            llm_funnel: None,
190            analyzers_run: vec![],
191            analyzers_skipped: vec![],
192        }
193    }
194
195    pub fn ran(&self) -> bool {
196        self.outcome == RunOutcome::Ran
197    }
198}
199
200/// The engine holds the registered analyzers, the host policy, and an optional
201/// LLM enrichment backend (§9).
202pub struct Engine {
203    analyzers: Vec<Box<dyn Analyzer>>,
204    policy: crate::policy::Policy,
205    /// Optional LLM backend. `None` → the DISCOVER/ENRICH stages are the
206    /// identity, so the pipeline is byte-for-byte the deterministic path.
207    llm: Option<Box<dyn crate::llm::LlmBackend>>,
208    /// Optional separate backend for the GROUND stage (§5.2, §11). `None` →
209    /// grounding rides `llm`. Lets a team point entailment at a cheaper or
210    /// specialized model (or take the generative model out of grounding
211    /// entirely) without changing the proposer/verifier.
212    ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
213}
214
215struct AnalysisPass {
216    survivors: Vec<Recommendation>,
217    proposed: u64,
218    deduped: u64,
219    analyzers_run: Vec<String>,
220    analyzers_skipped: Vec<AnalyzerSkip>,
221    llm_funnel: Option<LlmFunnel>,
222}
223
224impl Engine {
225    /// An engine with the default built-ins and a default (fully closed)
226    /// policy — nothing auto-applies, no LLM.
227    pub fn with_builtins() -> Self {
228        Engine {
229            analyzers: crate::analyzer::builtin_analyzers(),
230            policy: crate::policy::Policy::default(),
231            llm: None,
232            ground_llm: None,
233        }
234    }
235
236    /// An engine with no analyzers (register your own).
237    pub fn empty() -> Self {
238        Engine {
239            analyzers: vec![],
240            policy: crate::policy::Policy::default(),
241            llm: None,
242            ground_llm: None,
243        }
244    }
245
246    /// Install a host policy (the only place auto-apply is granted).
247    pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
248        self.policy = policy;
249        self
250    }
251
252    /// Attach an optional LLM enrichment backend (§9). Only ever *adds* cited
253    /// draft recommendations (stamped `origin = llm`, never auto-applied) and
254    /// whitelisted guidance notes — it can never gate or rewrite deterministic
255    /// output.
256    pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
257        self.llm = Some(backend);
258        self
259    }
260
261    /// Attach a separate backend for the GROUND stage (§5.2). Without this,
262    /// grounding uses the `with_llm` backend. Independent of the proposer so an
263    /// operator can run entailment on a cheaper/specialized model.
264    pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
265        self.ground_llm = Some(backend);
266        self
267    }
268
269    pub fn policy(&self) -> &crate::policy::Policy {
270        &self.policy
271    }
272
273    /// Register an additional analyzer (the linked-Rust seam).
274    pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
275        self.analyzers.push(analyzer);
276    }
277
278    pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
279        &self.analyzers
280    }
281
282    /// Run the exact production analysis/validation path without Phase 0
283    /// measurement or Phase 3 persistence. The immutable substrate borrow is
284    /// the replay safety boundary: recommendations, audit grains, state,
285    /// cooldowns, outcomes, and the op-log cannot be changed here.
286    ///
287    /// `overrides` is keyed by full analyzer id and overlays the file's stored
288    /// parameter map. Unknown keys fail closed through `resolve_params` and
289    /// surface as an analyzer skip, exactly as in a production run.
290    pub fn analyze_only<S: OmsSubstrate>(
291        &self,
292        sub: &S,
293        opts: &RunOptions,
294        overrides: &BTreeMap<String, Map<String, Value>>,
295        now_ms: i64,
296    ) -> Result<Vec<Recommendation>> {
297        let persisted = LoopPersisted::from_value(sub.load_state()?)?;
298        let analysis_watermark = if opts.full_sweep {
299            None
300        } else {
301            persisted.state.watermark_ms
302        };
303        Ok(self
304            .analysis_pass(
305                sub,
306                &persisted,
307                opts,
308                overrides,
309                analysis_watermark,
310                now_ms,
311                &[],
312            )?
313            .survivors)
314    }
315
316    /// Run one analysis pass. Idempotent under `dedup_key`; the watermark is
317    /// advanced at the end, so a crashed run simply re-runs.
318    pub fn run<S: OmsSubstrate>(
319        &self,
320        sub: &mut S,
321        opts: &RunOptions,
322        now_ms: i64,
323    ) -> Result<RunResult> {
324        let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
325        let watermark = persisted.state.watermark_ms;
326        // A full sweep analyzes the whole memory (watermark ignored for the
327        // analysis inputs), while gating, `new` counts, and the end-of-run
328        // watermark advance still use the real watermark. Dedup/cooldowns keep
329        // it from re-proposing what is already queued.
330        let analysis_watermark = if opts.full_sweep { None } else { watermark };
331
332        let (new_grains, new_error_events) = count_new(sub, watermark)?;
333        if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
334            return Ok(RunResult::skipped(reason, new_grains, new_error_events));
335        }
336
337        // Phase 0: re-measure applied recommendations due for review (the
338        // Verify gate). Records a measured outcome per due recommendation.
339        let outcome_inputs = measure_outcomes(sub, &mut persisted, now_ms)?;
340
341        let AnalysisPass {
342            survivors,
343            proposed,
344            deduped,
345            analyzers_run,
346            analyzers_skipped,
347            llm_funnel,
348        } = self.analysis_pass(
349            &*sub,
350            &persisted,
351            opts,
352            &BTreeMap::new(),
353            analysis_watermark,
354            now_ms,
355            &outcome_inputs,
356        )?;
357
358        // Phase 3 (needs &mut): store survivors + propose audit, then
359        // auto-apply the ones the host policy grants (all gates in §6.3).
360        let mut stored = 0u64;
361        let mut auto_applied = 0u64;
362        for mut rec in survivors {
363            let spec = rec.to_grain_spec(LOOP_NS)?;
364            let hash = sub.put_grain(&spec)?;
365            rec.hash = hash.clone();
366            let actor = format!("engine:{}", rec.analyzer);
367            let audit = AuditRecord {
368                rec_hash: hash.clone(),
369                from: None,
370                to: RecStatus::Pending,
371                actor: actor.clone(),
372                observer_type: ObserverType::System,
373                because: "analyzer proposed".into(),
374                previous_audit_hash: None,
375                gating: None,
376                at_ms: now_ms,
377            };
378            let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
379            persisted
380                .status_index
381                .insert(hash.clone(), RecStatus::Pending);
382            persisted.creators.insert(hash.clone(), actor);
383            // An LLM or external-command finding exists because someone ran
384            // it — record that principal too, so review can refuse the
385            // trigger approving their own model's output. Builtin analyzers
386            // stay engine-only: deterministic output has no human author.
387            if !matches!(rec.origin, Origin::Builtin) {
388                if let Some(trigger) = &opts.triggering_actor {
389                    persisted.co_creators.insert(hash.clone(), trigger.clone());
390                }
391            }
392            persisted.audit_heads.insert(hash.clone(), audit_hash);
393            stored += 1;
394
395            if self.can_auto_apply(&*sub, &rec) {
396                self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
397                auto_applied += 1;
398            }
399        }
400
401        persisted.state.last_run_ms = Some(now_ms);
402        persisted.state.watermark_ms = Some(now_ms);
403        sub.store_state(&persisted.to_value()?)?;
404
405        Ok(RunResult {
406            outcome: RunOutcome::Ran,
407            skip_reason: None,
408            new_grains,
409            new_error_events,
410            proposed,
411            deduped,
412            stored,
413            auto_applied,
414            analyzers_run,
415            analyzers_skipped,
416            llm_funnel,
417        })
418    }
419
420    /// Shared Phase 1–2 implementation for production and replay. Cooldowns
421    /// and the live recommendation queue are honored in both modes; only the
422    /// production caller proceeds into the mutating Phase 3 below.
423    #[allow(clippy::too_many_arguments)]
424    fn analysis_pass<S: OmsSubstrate>(
425        &self,
426        sub: &S,
427        persisted: &LoopPersisted,
428        opts: &RunOptions,
429        external_overrides: &BTreeMap<String, Map<String, Value>>,
430        analysis_watermark: Option<i64>,
431        now_ms: i64,
432        outcome_inputs: &[OutcomeInput],
433    ) -> Result<AnalysisPass> {
434        let existing = existing_dedup_keys(sub, persisted)?;
435        let mut analyzers_run = Vec::new();
436        let mut analyzers_skipped = Vec::new();
437        let mut candidates: Vec<Recommendation> = Vec::new();
438        let caps = sub.capabilities();
439
440        for analyzer in &self.analyzers {
441            let m = analyzer.manifest();
442            let cfg = persisted.config.get(&m.id);
443            let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
444            if !enabled {
445                analyzers_skipped.push(AnalyzerSkip {
446                    id: m.id.clone(),
447                    reason: "disabled".into(),
448                });
449                continue;
450            }
451            if self.policy.denies(m.family()) {
452                analyzers_skipped.push(AnalyzerSkip {
453                    id: m.id.clone(),
454                    reason: "denied by host policy".into(),
455                });
456                continue;
457            }
458            if let Some(missing) = missing_capability(m, caps) {
459                analyzers_skipped.push(AnalyzerSkip {
460                    id: m.id.clone(),
461                    reason: format!("missing capability: {missing}"),
462                });
463                continue;
464            }
465            let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
466            if let Some(extra) = external_overrides.get(&m.id) {
467                for (key, value) in extra {
468                    param_overrides.insert(key.clone(), value.clone());
469                }
470            }
471            let params = match m.resolve_params(&param_overrides) {
472                Ok(p) => p,
473                Err(e) => {
474                    analyzers_skipped.push(AnalyzerSkip {
475                        id: m.id.clone(),
476                        reason: e.to_string(),
477                    });
478                    continue;
479                }
480            };
481            let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
482            let ns_slice: &[String] = if ns_owned.is_empty() {
483                &opts.namespaces
484            } else {
485                &ns_owned
486            };
487            let reader: &dyn SubstrateRead = sub;
488            let ctx = AnalyzeCtx::new(
489                reader,
490                &params,
491                ns_slice,
492                analysis_watermark,
493                now_ms,
494                outcome_inputs,
495            );
496            match analyzer.analyze(&ctx) {
497                Ok(drafts) => {
498                    analyzers_run.push(m.id.clone());
499                    for draft in drafts {
500                        match stamp(m, &params, draft, now_ms) {
501                            Ok(rec) => candidates.push(rec),
502                            Err(e) => analyzers_skipped.push(AnalyzerSkip {
503                                id: m.id.clone(),
504                                reason: e.to_string(),
505                            }),
506                        }
507                    }
508                }
509                Err(e) => analyzers_skipped.push(AnalyzerSkip {
510                    id: m.id.clone(),
511                    reason: e.to_string(),
512                }),
513            }
514        }
515
516        let mut funnel = LlmFunnel::default();
517        if self.llm.is_some() {
518            candidates.extend(self.discover(
519                sub,
520                &candidates,
521                analysis_watermark,
522                &opts.namespaces,
523                now_ms,
524                &mut funnel,
525            ));
526        }
527
528        let proposed = candidates.len() as u64;
529        let mut seen = BTreeSet::new();
530        let mut survivors = Vec::new();
531        for candidate in candidates {
532            let family = crate::manifest::analyzer_family(&candidate.analyzer);
533            let floor = [
534                severity_floor_for(persisted, &candidate.analyzer),
535                self.policy.severity_floor(family),
536            ]
537            .into_iter()
538            .flatten()
539            .max();
540            if floor.is_some_and(|floor| candidate.severity < floor) {
541                continue;
542            }
543            if !seen.insert(candidate.dedup_key.clone()) {
544                continue;
545            }
546            if existing.contains(&candidate.dedup_key) {
547                continue;
548            }
549            if persisted
550                .cooldowns
551                .get(&candidate.dedup_key)
552                .is_some_and(|until| now_ms < *until)
553            {
554                continue;
555            }
556            survivors.push(candidate);
557        }
558        let deduped = proposed - survivors.len() as u64;
559        if self.llm.is_some() {
560            self.enrich(&mut survivors);
561        }
562        Ok(AnalysisPass {
563            survivors,
564            proposed,
565            deduped,
566            analyzers_run,
567            analyzers_skipped,
568            llm_funnel: self.llm.is_some().then_some(funnel),
569        })
570    }
571
572    /// DISCOVER (§9): ask the LLM for additional draft recommendations, given
573    /// the deterministic findings as *context* and a bounded, provenance-tagged
574    /// evidence bundle. Every returned draft must cite evidence present in the
575    /// bundle and target a memory/query surface; it is stamped `origin = llm`
576    /// (so it can never auto-apply) and enters the ordinary dedup/store path. A
577    /// failed or garbled response yields no drafts — never a failed run.
578    fn discover<S: OmsSubstrate>(
579        &self,
580        sub: &S,
581        candidates: &[Recommendation],
582        watermark: Option<i64>,
583        namespaces: &[String],
584        now_ms: i64,
585        funnel: &mut LlmFunnel,
586    ) -> Vec<Recommendation> {
587        let Some(llm) = &self.llm else {
588            return Vec::new();
589        };
590        let findings: Vec<crate::llm::FindingBrief> = candidates
591            .iter()
592            .take(32)
593            .map(|c| crate::llm::FindingBrief {
594                analyzer: c.analyzer.clone(),
595                summary: c.summary.render(),
596                target: c.target_ref.clone(),
597                severity: c.severity.as_str().to_string(),
598            })
599            .collect();
600        // Evidence bundle. Seeded first from the grains the deterministic
601        // findings cite, THEN — the non-parasitic step (§11) — topped up with
602        // RECENT grains (created since the last run) so the LLM gets its own
603        // lens and can find issues in grains no analyzer flagged. Without this
604        // the LLM could only elaborate near what determinism already caught.
605        let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
606        let mut bundle: BTreeSet<String> = BTreeSet::new();
607        let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
608        'cited: for c in candidates {
609            for h in &c.evidence {
610                // Every source gets a RESERVED share of the bundle, because a
611                // cap that only one source respects is not a budget. A single
612                // tool_failure finding may cite up to MAX_EVIDENCE (64)
613                // grains — the whole bundle — so without this the "give the
614                // LLM its own lens" seeding below could be starved to nothing
615                // by the very determinism it is supposed to look past. Found
616                // live: with the deterministic findings citing enough, the
617                // model never saw a single non-cited grain.
618                if evidence.len() >= CITED_SEED_CAP {
619                    break 'cited;
620                }
621                if !bundle.contains(h) {
622                    if let Ok(Some(g)) = sub.grain(h) {
623                        push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
624                    }
625                }
626            }
627        }
628        let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
629            vec![None]
630        } else {
631            namespaces.iter().map(|n| Some(n.as_str())).collect()
632        };
633        let opts = ReadOpts { live_only: true, since_ms: watermark };
634        // Tool grains carry the raw experience of a tool-using agent, and an
635        // ERROR is the part a reflection pass can act on. Seeded after the
636        // cited grains and inside its own reserved share, so a busy desk
637        // cannot crowd out the facts and observations below.
638        //
639        // Without this the LLM saw tool failures only through the
640        // deterministic findings that happened to cite them: it could
641        // elaborate on what clustering already caught, but could never find
642        // a failure clustering missed — the one thing it is here for. The
643        // top-up below called itself "non-parasitic" while omitting the very
644        // grain type the flagship analyzer reads.
645        let mut tool_seeded = 0usize;
646        'tools: for ns in &scan_ns {
647            if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
648                for g in recent {
649                    if tool_seeded >= TOOL_SEED_CAP
650                        || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
651                    {
652                        break 'tools;
653                    }
654                    if !g.is_error() {
655                        continue;
656                    }
657                    let before = evidence.len();
658                    push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
659                    if evidence.len() > before {
660                        tool_seeded += 1;
661                    }
662                }
663            }
664        }
665        // Observations BEFORE facts, with their own small reserve.
666        //
667        // An Observation is where a human's own words land — a supervisor's
668        // note, a reviewer's correction, an instruction for next time. That
669        // signal is stated ONCE by nature, so recency and frequency seeding
670        // structurally bury it: one note loses to three hundred routine
671        // records every time, and the rarest evidence is usually the most
672        // valuable. Exhausting facts first (as this did) meant a desk with
673        // any volume showed the model no human input at all.
674        'notes: for ns in &scan_ns {
675            if let Ok(recent) =
676                sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
677            {
678                for g in recent {
679                    if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
680                        break 'notes;
681                    }
682                    push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
683                }
684            }
685        }
686        'seed: for gt in [
687            crate::model::grain_type::FACT,
688            crate::model::grain_type::OBSERVATION,
689        ] {
690            for ns in &scan_ns {
691                if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
692                    for g in recent {
693                        if evidence.len() >= EVIDENCE_CAP {
694                            break 'seed;
695                        }
696                        push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g);
697                    }
698                }
699            }
700        }
701        funnel.evidence = evidence.len() as u64;
702        if evidence.is_empty() {
703            return Vec::new(); // nothing to reflect on
704        }
705        // PROPOSE (§5.1): the abstention-legitimate objective — "nothing to
706        // report" is a first-class, zero-penalty answer. The operator-taste
707        // history (recent approve/reject decisions on llm findings) is passed so
708        // the model learns what this reviewer accepts.
709        let (approved, rejected) = self.llm_history(sub);
710        let request = crate::llm::LlmRequest {
711            loop_proto: 1,
712            op: "discover",
713            instructions: DISCOVER_INSTRUCTIONS,
714            findings: findings.clone(),
715            evidence: evidence.clone(),
716            rejected,
717            approved,
718        };
719        let Ok(body) = serde_json::to_string(&request) else {
720            return Vec::new();
721        };
722        let raw = match llm.complete(&body) {
723            Ok(r) => r,
724            Err(_) => return Vec::new(), // fail-soft
725        };
726        // Cheap structural validation (cite-check + target class); collect the
727        // survivors for the verifier. Storing the normalized target string
728        // avoids a TargetRef clone through the pipeline.
729        let caps = sub.capabilities();
730        let mut validated: Vec<ValidatedDraft> = Vec::new();
731        let drafts: Vec<_> = crate::llm::parse_discover(&raw)
732            .recommendations
733            .into_iter()
734            .take(crate::llm::MAX_LLM_DRAFTS)
735            .collect();
736        funnel.proposed = drafts.len() as u64;
737        for d in drafts {
738            let cited: Vec<String> =
739                d.evidence.iter().filter(|h| bundle.contains(*h)).cloned().collect();
740            if cited.is_empty() {
741                funnel.dropped_uncited += 1;
742                continue; // uncited → drop (no fabrication)
743            }
744            let Ok(target) = TargetRef::parse(&d.target) else {
745                funnel.dropped_target += 1;
746                continue;
747            };
748            let tc = target.target_class();
749            // The classes the proposal vocabulary can reach: memory
750            // (entity/grain), query (query/template) and code (tool). The
751            // prompt (`doc:`), `host:`, `evalset:` and `model:` classes stay
752            // closed to the model — it may not rewrite the agent's prompt, its
753            // host config, or the gate that grades its own code.
754            if !matches!(tc, "memory" | "query" | "code") {
755                funnel.dropped_target += 1;
756                continue;
757            }
758            // Resolve BEFORE the gates: what GROUND entails and VERIFY
759            // stress-tests is exactly what an apply would do.
760            let resolved = resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps);
761            // A `tool:` target has exactly one legal shape (Rule E1: a code
762            // target REQUIRES action_kind code_revision), so an unresolved one
763            // could not even be stamped advisory — drop it here rather than
764            // spend two model calls on something that fails validation after.
765            if tc == "code" && resolved.is_none() {
766                funnel.dropped_target += 1;
767                continue;
768            }
769            validated.push(ValidatedDraft {
770                draft: d,
771                target_ref: target.as_string(),
772                cited,
773                resolved,
774            });
775        }
776        funnel.cited = validated.len() as u64;
777        if validated.is_empty() {
778            return Vec::new();
779        }
780        // GROUND → VERIFY → ROUTE (§5.2–5.4): only drafts that survive an
781        // independent grounding entailment check *and* an adversarial
782        // verification pass (each a separate call — proposer ≠ scorer) reach the
783        // queue, stamped with the verifier's calibrated confidence.
784        // GROUND may run on a separate backend (§11); VERIFY always uses the
785        // main llm (the proposer≠scorer independence is on VERIFY, not GROUND).
786        let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
787        self.verify_drafts(&**llm, ground, validated, &evidence, now_ms, funnel)
788    }
789
790    /// GROUND → VERIFY → ROUTE (§5.2–5.4). Two independent model calls, batched
791    /// over the drafts: a grounding-entailment gate ("does the cited evidence
792    /// support the claim?"), then an adversarial keep/kill with a calibrated
793    /// confidence. A draft reaches the queue only if it is grounded **and** kept
794    /// **and** clears the confidence floor. Any failed call drops the whole LLM
795    /// contribution for the run (safe default), never the run.
796    fn verify_drafts(
797        &self,
798        llm: &dyn crate::llm::LlmBackend,
799        ground: &dyn crate::llm::LlmBackend,
800        validated: Vec<ValidatedDraft>,
801        evidence: &[crate::llm::EvidenceItem],
802        now_ms: i64,
803        funnel: &mut LlmFunnel,
804    ) -> Vec<Recommendation> {
805        use crate::llm::*;
806        let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
807            evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
808        let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
809            cited
810                .iter()
811                .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
812                .collect()
813        };
814
815        // GROUND (§5.2): decompose-then-entail per draft, batched into one call.
816        // The claim includes any authored lesson — the gate must judge exactly
817        // what an apply would record, not only the finding's summary.
818        let claims: Vec<GroundItem> = validated
819            .iter()
820            .enumerate()
821            .map(|(i, v)| GroundItem {
822                id: i,
823                claim: claim_text(&v.draft, v.resolved.as_ref()),
824                evidence: ev_for(&v.cited),
825            })
826            .collect();
827        let ground_req = GroundRequest {
828            loop_proto: 1,
829            op: "ground",
830            instructions: GROUND_INSTRUCTIONS,
831            claims,
832        };
833        // A grounding pass that REFUSED every draft and one that never
834        // answered are the same number of survivors and opposite problems:
835        // the first is the gate doing its job, the second is a backend having
836        // a bad minute while the engine fail-softs. Count the verdicts
837        // actually returned so the two are distinguishable afterwards.
838        let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
839            .ok()
840            .and_then(|b| ground.complete(&b).ok())
841        {
842            Some(raw) => {
843                let parsed = parse_ground(&raw);
844                funnel.ground_verdicts = parsed.results.len() as u64;
845                parsed
846                    .results
847                    .into_iter()
848                    .filter(|r| r.supported)
849                    .map(|r| r.id)
850                    .collect()
851            }
852            None => {
853                funnel.ground_call_failed = true;
854                return Vec::new();
855            }
856        };
857        funnel.grounded = grounded.len() as u64;
858        if grounded.is_empty() {
859            return Vec::new();
860        }
861
862        // VERIFY (§5.3): adversarial keep/kill over the grounded drafts, a
863        // separate call from the proposer. Soundness + abstention only — NOT
864        // novelty. Novelty is steered at DISCOVER and settled by human review;
865        // asking a weak verifier to judge it just makes it hallucinate "already
866        // known" and kill genuine findings (§11).
867        let items: Vec<VerifyItem> = validated
868            .iter()
869            .enumerate()
870            .filter(|(i, _)| grounded.contains(i))
871            .map(|(i, v)| VerifyItem {
872                id: i,
873                // Same rule as GROUND: the adversarial pass sees the change.
874                summary: claim_text(&v.draft, v.resolved.as_ref()),
875                target: v.target_ref.clone(),
876                evidence: ev_for(&v.cited),
877            })
878            .collect();
879        let verify_req = VerifyRequest {
880            loop_proto: 1,
881            op: "verify",
882            instructions: VERIFY_INSTRUCTIONS,
883            findings: items,
884        };
885        let verdicts: std::collections::BTreeMap<usize, f64> =
886            match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
887                Some(raw) => parse_verify(&raw)
888                    .results
889                    .into_iter()
890                    .filter(|r| r.keep)
891                    .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
892                    .collect(),
893                None => return Vec::new(),
894            };
895
896        funnel.kept = verdicts.len() as u64;
897        // ROUTE (§5.4): grounded ∧ kept ∧ verifier-confidence ≥ floor. The
898        // verifier's confidence (the independent signal) is what we trust and
899        // stamp — not the proposer's self-report.
900        let mut out = Vec::new();
901        for (i, v) in validated.into_iter().enumerate() {
902            if let Some(&conf) = verdicts.get(&i) {
903                if conf >= MIN_LLM_CONFIDENCE {
904                    out.push(stamp_llm(
905                        llm.model(),
906                        &v.draft,
907                        v.target_ref,
908                        v.cited,
909                        v.resolved,
910                        conf,
911                        now_ms,
912                    ));
913                }
914            }
915        }
916        funnel.stored = out.len() as u64;
917        out
918    }
919
920    /// Recent operator decisions on `origin = llm` findings — approved (incl.
921    /// applied) and rejected summaries, most-recent first and bounded — so
922    /// DISCOVER can learn what this reviewer accepts (§9). Best-effort: a read
923    /// failure yields empty history, never an error.
924    fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
925        const MAX: usize = 20;
926        let Ok(mut recs) = self.recommendations(sub, None) else {
927            return (Vec::new(), Vec::new());
928        };
929        recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
930        recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
931        let mut approved = Vec::new();
932        let mut rejected = Vec::new();
933        for r in &recs {
934            match r.status {
935                RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
936                    if approved.len() < MAX =>
937                {
938                    approved.push(r.summary.render());
939                }
940                RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
941                _ => {}
942            }
943        }
944        (approved, rejected)
945    }
946
947    /// ENRICH (§9): ask the LLM to add a short guidance note to the surviving
948    /// deterministic recommendations. Whitelist-only — only `guidance` is
949    /// merged (capped), and only onto recs that don't already have one; the
950    /// engine-templated summary is never touched. Fail-soft.
951    fn enrich(&self, survivors: &mut [Recommendation]) {
952        let Some(llm) = &self.llm else {
953            return;
954        };
955        if survivors.is_empty() {
956            return;
957        }
958        let findings: Vec<crate::llm::FindingBrief> = survivors
959            .iter()
960            .map(|r| crate::llm::FindingBrief {
961                analyzer: r.analyzer.clone(),
962                summary: r.summary.render(),
963                target: r.target_ref.clone(),
964                severity: r.severity.as_str().to_string(),
965            })
966            .collect();
967        let request = crate::llm::LlmRequest {
968            loop_proto: 1,
969            op: "enrich",
970            instructions: ENRICH_INSTRUCTIONS,
971            findings,
972            evidence: Vec::new(),
973            rejected: Vec::new(),
974            approved: Vec::new(),
975        };
976        let Ok(body) = serde_json::to_string(&request) else {
977            return;
978        };
979        let raw = match llm.complete(&body) {
980            Ok(r) => r,
981            Err(_) => return,
982        };
983        for note in crate::llm::parse_enrich(&raw).notes {
984            if note.guidance.trim().is_empty() {
985                continue;
986            }
987            if let Some(r) = survivors
988                .iter_mut()
989                .find(|r| r.target_ref == note.target && r.guidance.is_none())
990            {
991                r.guidance = Some(crate::llm::cap(&note.guidance, crate::llm::MAX_GUIDANCE_LEN));
992            }
993        }
994    }
995
996    /// Evaluate the auto-apply gate (§6.3) — ALL preconditions must hold:
997    /// host opt-in + policy grant, builtin origin, memory/query target,
998    /// non-destructive, and engine-side shape verification: SUPERSEDE-only
999    /// structural curation (never an ADD that introduces evidence-derived
1000    /// text) whose every replacement is **value-identical** to the grain it
1001    /// supersedes (the exact-equality check — a near-duplicate consolidation
1002    /// stays pending). A default (closed) policy never grants, so nothing
1003    /// auto-applies.
1004    fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
1005        if !rec.origin.auto_apply_eligible() || rec.destructive {
1006            return false;
1007        }
1008        // The analyzer must declare its curation auto-appliable. An analyzer
1009        // whose manifest is `Never` (e.g. fork surfacing — a lossy merge) is
1010        // never auto-applied even if the payload passes the shape check.
1011        let manifest_ok = self
1012            .analyzers
1013            .iter()
1014            .map(|a| a.manifest())
1015            .find(|m| m.id == rec.analyzer)
1016            .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
1017        if !manifest_ok {
1018            return false;
1019        }
1020        let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1021            return false;
1022        };
1023        let family = crate::manifest::analyzer_family(&rec.analyzer);
1024        if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1025            return false;
1026        }
1027        // Shape verification: only a CAL batch of pure SUPERSEDE statements
1028        // whose replacements change no value is structural curation. An ADD
1029        // (introducing content), a FORGET (destructive), or a supersession
1030        // that alters any field disqualifies.
1031        match &rec.proposal {
1032            Proposal::Cal { cal } => cal
1033                .lines()
1034                .map(str::trim)
1035                .filter(|l| !l.is_empty())
1036                .all(|l| supersede_is_value_identical(sub, l)),
1037            _ => false,
1038        }
1039    }
1040
1041    /// Apply a recommendation as `policy:auto` (the only `pending → applied`
1042    /// path). Records the applied inverse + a hash-chained audit grain.
1043    fn auto_apply<S: OmsSubstrate>(
1044        &self,
1045        sub: &mut S,
1046        p: &mut LoopPersisted,
1047        rec: &Recommendation,
1048        now_ms: i64,
1049    ) -> Result<()> {
1050        let mut created = Vec::new();
1051        if let Proposal::Cal { cal } = &rec.proposal {
1052            // Belt and braces over the policy: `grants_auto_apply` already
1053            // excludes the `query` class, so a definition rewrite cannot reach
1054            // this path. If one ever did, it would apply with no recorded
1055            // inverse and no human BECAUSE — refuse instead.
1056            if cal.lines().map(str::trim).any(is_definition_statement) {
1057                return Err(Error::InvalidProposal(
1058                    "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1059                     auto-applied: it changes what every future context contains, so it \
1060                     requires a human APPROVE + APPLY with BECAUSE"
1061                        .into(),
1062                ));
1063            }
1064            for r in sub.execute_cal(cal)? {
1065                if let Some(h) = r.get("hash").and_then(Value::as_str) {
1066                    created.push(h.to_string());
1067                }
1068            }
1069        }
1070        let applied = AppliedRecord {
1071            applied_at_ms: now_ms,
1072            target_ref: rec.target_ref.clone(),
1073            rollbackable: rec.rollbackable,
1074            created_hashes: created,
1075            inverse_cal: None,
1076            metric: rec.metric.clone(),
1077        };
1078        let prev = p.audit_heads.get(&rec.hash).cloned();
1079        let audit = AuditRecord {
1080            rec_hash: rec.hash.clone(),
1081            from: Some(RecStatus::Pending),
1082            to: RecStatus::Applied,
1083            actor: "policy:auto".into(),
1084            observer_type: ObserverType::Policy,
1085            because: "auto-applied per host policy".into(),
1086            previous_audit_hash: prev,
1087            gating: None,
1088            at_ms: now_ms,
1089        };
1090        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1091        p.audit_heads.insert(rec.hash.clone(), audit_hash);
1092        p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1093        p.applied.insert(rec.hash.clone(), applied);
1094        Ok(())
1095    }
1096
1097    /// Approve or reject a pending recommendation. Requires the `review` scope,
1098    /// a mandatory BECAUSE, and blocks self-approval against the creating actor.
1099    #[allow(clippy::too_many_arguments)]
1100    pub fn review<S: OmsSubstrate>(
1101        &self,
1102        sub: &mut S,
1103        rec_hash: &str,
1104        decision: Decision,
1105        actor: &str,
1106        observer: ObserverType,
1107        scopes: &ScopeSet,
1108        because: &str,
1109        now_ms: i64,
1110    ) -> Result<()> {
1111        if !scopes.has(Scope::Review) {
1112            return Err(Error::ScopeDenied("review".into()));
1113        }
1114        let because = validate_because(because)?;
1115        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1116        let status = *p
1117            .status_index
1118            .get(rec_hash)
1119            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1120        let to = match decision {
1121            Decision::Approve => RecStatus::Approved,
1122            Decision::Reject => RecStatus::Rejected,
1123        };
1124        if !status.can_transition_to(to, false) {
1125            return Err(Error::LifecycleViolation(format!(
1126                "{} -> {}",
1127                status.as_str(),
1128                to.as_str()
1129            )));
1130        }
1131        if to == RecStatus::Approved {
1132            if let Some(creator) = p.creators.get(rec_hash) {
1133                if creator == actor {
1134                    return Err(Error::SelfApproval(format!(
1135                        "{actor} created this recommendation"
1136                    )));
1137                }
1138            }
1139            if let Some(trigger) = p.co_creators.get(rec_hash) {
1140                if trigger == actor {
1141                    return Err(Error::SelfApproval(format!(
1142                        "{actor} triggered the run that authored this recommendation"
1143                    )));
1144                }
1145            }
1146        }
1147        let prev = p.audit_heads.get(rec_hash).cloned();
1148        let audit = AuditRecord {
1149            rec_hash: rec_hash.into(),
1150            from: Some(status),
1151            to,
1152            actor: actor.into(),
1153            observer_type: observer,
1154            because,
1155            previous_audit_hash: prev,
1156            gating: None,
1157            at_ms: now_ms,
1158        };
1159        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1160        p.audit_heads.insert(rec_hash.into(), audit_hash);
1161        p.status_index.insert(rec_hash.into(), to);
1162        if to == RecStatus::Rejected {
1163            if let Ok(rec) = load_rec(sub, rec_hash) {
1164                // Exponential backoff keyed on dedup_key: 7d, 14d, 28d, … capped
1165                // at 90d, so a finding a reviewer keeps rejecting stops
1166                // re-surfacing on a fixed 7d cadence (was a flat 7d despite the
1167                // "doubling" comment).
1168                const BASE_MS: i64 = 7 * 86_400_000;
1169                const CAP_MS: i64 = 90 * 86_400_000;
1170                let strikes = p.cooldown_strikes.entry(rec.dedup_key.clone()).or_insert(0);
1171                let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
1172                *strikes = strikes.saturating_add(1);
1173                p.cooldowns.insert(rec.dedup_key, now_ms + interval);
1174            }
1175        }
1176        sub.store_state(&p.to_value()?)?;
1177        Ok(())
1178    }
1179
1180    /// Check everything [`apply`](Self::apply) would refuse on, without writing
1181    /// anything.
1182    ///
1183    /// This exists for the fused approve-and-apply callers (the bindings'
1184    /// `apply_recommendation`). Recording the approval first and *then* hitting
1185    /// the destructive gate strands the recommendation in `approved`, which has
1186    /// no exit but `applied` or `expired` — `approved → rejected` is not a
1187    /// legal transition — so a refused apply left the reviewer unable to
1188    /// dismiss it. Ask first, then approve.
1189    ///
1190    /// Deliberately does not check the lifecycle transition: the caller is
1191    /// about to make it legal by approving.
1192    /// `has_gating` is whether the caller will supply a gating run at apply:
1193    /// a gated revision (code or adapter) without one is refused HERE, before
1194    /// a fused approve-and-apply records the approval — `approved` has no
1195    /// exit but `applied` or `expired`, so asking after would strand it.
1196    pub fn preflight_apply<S: OmsSubstrate>(
1197        &self,
1198        sub: &S,
1199        rec_hash: &str,
1200        scopes: &ScopeSet,
1201        allow_destructive: bool,
1202        has_gating: bool,
1203    ) -> Result<()> {
1204        if !scopes.has(Scope::Apply) {
1205            return Err(Error::ScopeDenied("apply".into()));
1206        }
1207        let rec = load_rec(sub, rec_hash)?;
1208        if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1209            return Err(Error::DestructiveGated(
1210                "destructive apply requires admin scope + allow_destructive".into(),
1211            ));
1212        }
1213        ensure_executable(rec.action_kind, &rec.proposal)?;
1214        if requires_gating(rec.action_kind) && !has_gating {
1215            return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1216        }
1217        Ok(())
1218    }
1219
1220    /// Apply an approved recommendation. Requires `apply`; destructive payloads
1221    /// additionally require `admin` + `allow_destructive`. Records the applied
1222    /// info (inverse plan) for rollback.
1223    #[allow(clippy::too_many_arguments)]
1224    pub fn apply<S: OmsSubstrate>(
1225        &self,
1226        sub: &mut S,
1227        rec_hash: &str,
1228        actor: &str,
1229        observer: ObserverType,
1230        scopes: &ScopeSet,
1231        because: &str,
1232        allow_destructive: bool,
1233        now_ms: i64,
1234    ) -> Result<AppliedRecord> {
1235        self.apply_inner(
1236            sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1237        )
1238    }
1239
1240    /// Load the gating evidence for `rec_hash` from the RECORDED
1241    /// `mg:eval_run` summary named by `run_id` — the one loader every
1242    /// surface (CLI, bindings, MCP, HTTP) shares, so the stats that admit a
1243    /// gated revision can never come from a caller: they are read back from
1244    /// the Fact `areev eval run` journaled, within the recommendation's own
1245    /// pinned evalset (a run id from a different evalset simply isn't found).
1246    pub fn gating_evidence<S: OmsSubstrate>(
1247        &self,
1248        sub: &S,
1249        rec_hash: &str,
1250        run_id: &str,
1251    ) -> Result<crate::recommendation::GatingEvidence> {
1252        let rec = self
1253            .recommendations(sub, None)?
1254            .into_iter()
1255            .find(|r| r.hash == rec_hash)
1256            .ok_or_else(|| {
1257                Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1258            })?;
1259        let pin = rec.evalset_hash.ok_or_else(|| {
1260            Error::InvalidProposal(
1261                "this recommendation pins no evalset — a gating run applies only \
1262                 to code and adapter revisions"
1263                    .into(),
1264            )
1265        })?;
1266        // Read through the shared evalset reader, so the run that GATES an
1267        // apply and the runs that later JUDGE it are parsed by exactly one
1268        // piece of code — two parsers drifting apart would let a rule be
1269        // admitted on one reading of a summary and measured on another.
1270        match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1271            Some(run) => Ok(crate::recommendation::GatingEvidence {
1272                evalset_hash: pin,
1273                run_id: run.run_id,
1274                passed: run.passed,
1275                failed: run.failed,
1276            }),
1277            None => Err(Error::InvalidProposal(format!(
1278                "no recorded gate run '{run_id}' for evalset {pin} — run \
1279                 `areev eval run --evalset {pin} ...` first"
1280            ))),
1281        }
1282    }
1283
1284    /// Apply WITH the §7.4 evalset-run edge — the only path that can apply a
1285    /// gated (code or adapter) revision. The evidence is validated against
1286    /// the recommendation's pin and recorded on the audit Observation.
1287    #[allow(clippy::too_many_arguments)]
1288    pub fn apply_gated<S: OmsSubstrate>(
1289        &self,
1290        sub: &mut S,
1291        rec_hash: &str,
1292        actor: &str,
1293        observer: ObserverType,
1294        scopes: &ScopeSet,
1295        because: &str,
1296        allow_destructive: bool,
1297        gating: &crate::recommendation::GatingEvidence,
1298        now_ms: i64,
1299    ) -> Result<AppliedRecord> {
1300        self.apply_inner(
1301            sub,
1302            rec_hash,
1303            actor,
1304            observer,
1305            scopes,
1306            because,
1307            allow_destructive,
1308            Some(gating),
1309            now_ms,
1310        )
1311    }
1312
1313    #[allow(clippy::too_many_arguments)]
1314    fn apply_inner<S: OmsSubstrate>(
1315        &self,
1316        sub: &mut S,
1317        rec_hash: &str,
1318        actor: &str,
1319        observer: ObserverType,
1320        scopes: &ScopeSet,
1321        because: &str,
1322        allow_destructive: bool,
1323        gating: Option<&crate::recommendation::GatingEvidence>,
1324        now_ms: i64,
1325    ) -> Result<AppliedRecord> {
1326        if !scopes.has(Scope::Apply) {
1327            return Err(Error::ScopeDenied("apply".into()));
1328        }
1329        let because = validate_because(because)?;
1330        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1331        let status = *p
1332            .status_index
1333            .get(rec_hash)
1334            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1335        if !status.can_transition_to(RecStatus::Applied, false) {
1336            return Err(Error::LifecycleViolation(format!(
1337                "{} -> applied (approve first)",
1338                status.as_str()
1339            )));
1340        }
1341        let rec = load_rec(sub, rec_hash)?;
1342        if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1343            return Err(Error::DestructiveGated(
1344                "destructive apply requires admin scope + allow_destructive".into(),
1345            ));
1346        }
1347        // §7.4: a code or adapter revision applies ONLY through the
1348        // evalset-run edge — the pin must match, the pinned evalset must
1349        // still be LIVE (a superseded evalset invalidates in-flight
1350        // recommendations: they re-gate), and a failing gate admits nothing.
1351        if requires_gating(rec.action_kind) {
1352            let g = gating
1353                .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1354            let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1355            if g.evalset_hash != pin {
1356                return Err(Error::InvalidProposal(format!(
1357                    "gating ran evalset {} but the recommendation is pinned \
1358                     to {pin} (Rule E1)",
1359                    g.evalset_hash
1360                )));
1361            }
1362            match sub.grain(pin)? {
1363                Some(evalset) if evalset.is_live() => {}
1364                Some(_) => {
1365                    return Err(Error::InvalidProposal(
1366                        "the pinned evalset was superseded after gating — \
1367                         the recommendation must re-gate (Rule E1)"
1368                            .into(),
1369                    ))
1370                }
1371                None => {
1372                    return Err(Error::InvalidProposal(format!(
1373                        "pinned evalset {pin} not found in the substrate"
1374                    )))
1375                }
1376            }
1377            if g.failed > 0 {
1378                return Err(Error::InvalidProposal(format!(
1379                    "the gating run failed {}/{} cases — a failing gate \
1380                     admits nothing",
1381                    g.failed,
1382                    g.passed + g.failed
1383                )));
1384            }
1385        }
1386
1387        // Execute the proposal.
1388        let mut created = Vec::new();
1389        // The inverse of a change that creates no grain — see
1390        // `AppliedRecord::inverse_cal`. Captured BEFORE execution, because
1391        // afterwards the previous definition is gone.
1392        let mut inverse_cal: Option<String> = None;
1393        match &rec.proposal {
1394            Proposal::Cal { cal } => {
1395                for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1396                    if !is_definition_statement(line) {
1397                        continue;
1398                    }
1399                    match sub.definition_inverse(line)? {
1400                        Some(inv) => inverse_cal = Some(inv),
1401                        None => {
1402                            return Err(Error::InvalidProposal(format!(
1403                                "this substrate cannot record a rollback inverse for {line:?}; \
1404                                 a definition rewrite that ROLLBACK could not undo is refused \
1405                                 rather than applied"
1406                            )))
1407                        }
1408                    }
1409                }
1410                let rows = sub.execute_cal(cal)?;
1411                for r in rows {
1412                    if let Some(h) = r.get("hash").and_then(Value::as_str) {
1413                        created.push(h.to_string());
1414                    }
1415                }
1416            }
1417            // The engine has no executable Edit primitive. Marking this
1418            // Applied used to be a lie (and rollback had no inverse).
1419            Proposal::Edit { .. } => {
1420                return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1421            }
1422            // A gated code or adapter revision EXECUTES by writing the
1423            // promotion grain: an immutable record that this target now
1424            // resolves to the proposed code (a §7.4 blob address) or adapter
1425            // (the tuning seam's registry tuple). Hosts read the promotion
1426            // to re-resolve — `(tool:X, mg:code_promotion)` or
1427            // `(model:X, mg:adapter_promotion)`; retracting it is the
1428            // rollback inverse, so the apply is rollbackable end-to-end.
1429            Proposal::Data { data } if requires_gating(rec.action_kind) => {
1430                let relation = if rec.action_kind == ActionKind::AdapterRevision {
1431                    "mg:adapter_promotion"
1432                } else {
1433                    "mg:code_promotion"
1434                };
1435                // A code revision authored by DISCOVER carries its SOURCE
1436                // inline, because the discovery pass reads the substrate and
1437                // cannot write to it. Move it into the CAS here so the
1438                // promotion names a content ADDRESS: §7.4's rule that code
1439                // enters the substrate only through the blob seam holds
1440                // however the revision was authored, and the promotion grain
1441                // stays a pointer rather than swelling to hold a program.
1442                let mut promoted = data.clone();
1443                if let Some(Value::String(src)) = promoted.remove("source") {
1444                    let address = sub.put_blob(src.as_bytes())?;
1445                    promoted.insert("code_address".into(), Value::from(address));
1446                }
1447                let mut spec = crate::substrate::GrainSpec::new(
1448                    crate::model::grain_type::FACT,
1449                    LOOP_NS,
1450                )
1451                .with_field("subject", rec.target_ref.clone())
1452                .with_field("relation", relation)
1453                .with_field(
1454                    "object",
1455                    serde_json::to_string(&promoted)
1456                        .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1457                )
1458                .with_field("rec_hash", rec_hash.to_string());
1459                if let Some(g) = gating {
1460                    spec = spec
1461                        .with_field("gating_evalset", g.evalset_hash.clone())
1462                        .with_field("gating_run_id", g.run_id.clone());
1463                }
1464                created.push(sub.put_grain(&spec)?);
1465            }
1466            Proposal::Data { data } => {
1467                // OutcomeReview is the one executable Data shape: its
1468                // `revert_of` points at an earlier applied recommendation.
1469                // Reuse the ordinary rollback path so the created hashes are
1470                // really retracted and the original lifecycle/audit advances.
1471                let revert_of = data
1472                    .get("revert_of")
1473                    .and_then(Value::as_str)
1474                    .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1475                self.rollback(
1476                    sub,
1477                    revert_of,
1478                    actor,
1479                    observer,
1480                    scopes,
1481                    &because,
1482                    now_ms,
1483                )?;
1484                // rollback stored a newer lifecycle state; merge this apply
1485                // into that state rather than overwriting the rollback.
1486                p = LoopPersisted::from_value(sub.load_state()?)?;
1487            }
1488        }
1489
1490        let applied = AppliedRecord {
1491            applied_at_ms: now_ms,
1492            target_ref: rec.target_ref.clone(),
1493            rollbackable: rec.rollbackable,
1494            created_hashes: created,
1495            inverse_cal,
1496            metric: rec.metric.clone(),
1497        };
1498        let prev = p.audit_heads.get(rec_hash).cloned();
1499        let audit = AuditRecord {
1500            rec_hash: rec_hash.into(),
1501            from: Some(status),
1502            to: RecStatus::Applied,
1503            actor: actor.into(),
1504            observer_type: observer,
1505            because,
1506            previous_audit_hash: prev,
1507            gating: gating.cloned(),
1508            at_ms: now_ms,
1509        };
1510        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1511        p.audit_heads.insert(rec_hash.into(), audit_hash);
1512        p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1513        p.applied.insert(rec_hash.into(), applied.clone());
1514        sub.store_state(&p.to_value()?)?;
1515        Ok(applied)
1516    }
1517
1518    /// Roll back an applied recommendation by retracting the grains it created.
1519    /// Fails for non-rollbackable applies (e.g. FORGET).
1520    #[allow(clippy::too_many_arguments)]
1521    pub fn rollback<S: OmsSubstrate>(
1522        &self,
1523        sub: &mut S,
1524        rec_hash: &str,
1525        actor: &str,
1526        observer: ObserverType,
1527        scopes: &ScopeSet,
1528        because: &str,
1529        now_ms: i64,
1530    ) -> Result<()> {
1531        if !scopes.has(Scope::Apply) {
1532            return Err(Error::ScopeDenied("apply".into()));
1533        }
1534        let because = validate_because(because)?;
1535        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1536        let status = *p
1537            .status_index
1538            .get(rec_hash)
1539            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1540        if !status.can_transition_to(RecStatus::RolledBack, false) {
1541            return Err(Error::LifecycleViolation(format!(
1542                "{} -> rolled_back",
1543                status.as_str()
1544            )));
1545        }
1546        let applied = p
1547            .applied
1548            .get(rec_hash)
1549            .cloned()
1550            .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1551        if !applied.rollbackable {
1552            return Err(Error::LifecycleViolation(
1553                "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1554            ));
1555        }
1556        for h in &applied.created_hashes {
1557            sub.retract(h, &format!("rollback of {rec_hash}"))?;
1558        }
1559        // A definition rewrite creates no grain, so retracting `created_hashes`
1560        // undoes nothing. Restoring it means re-running the statement captured
1561        // at apply time — the previous definition, or a DROP when there was
1562        // none. Runs BEFORE the audit is written, so a failed restore leaves
1563        // the recommendation `applied` (still true) rather than recording a
1564        // rollback that did not happen.
1565        if let Some(inverse) = &applied.inverse_cal {
1566            sub.execute_cal(inverse)?;
1567        }
1568        let prev = p.audit_heads.get(rec_hash).cloned();
1569        let audit = AuditRecord {
1570            rec_hash: rec_hash.into(),
1571            from: Some(status),
1572            to: RecStatus::RolledBack,
1573            actor: actor.into(),
1574            observer_type: observer,
1575            because,
1576            previous_audit_hash: prev,
1577            gating: None,
1578            at_ms: now_ms,
1579        };
1580        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1581        p.audit_heads.insert(rec_hash.into(), audit_hash);
1582        p.status_index
1583            .insert(rec_hash.into(), RecStatus::RolledBack);
1584        sub.store_state(&p.to_value()?)?;
1585        Ok(())
1586    }
1587
1588    /// List stored recommendations, optionally filtered by status. Status comes
1589    /// from the rebuildable index, not the immutable grain body. Ordered for
1590    /// review triage — highest severity first, then oldest first — and stable
1591    /// across runs for identical input.
1592    pub fn recommendations<S: OmsSubstrate>(
1593        &self,
1594        sub: &S,
1595        status_filter: Option<RecStatus>,
1596    ) -> Result<Vec<Recommendation>> {
1597        let p = LoopPersisted::from_value(sub.load_state()?)?;
1598        let grains = sub.grains_of_type(
1599            crate::model::grain_type::RECOMMENDATION,
1600            Some(LOOP_NS),
1601            ReadOpts {
1602                live_only: false,
1603                since_ms: None,
1604            },
1605        )?;
1606        let mut out = Vec::new();
1607        for g in grains {
1608            let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1609            rec.status = p
1610                .status_index
1611                .get(&g.hash)
1612                .copied()
1613                .unwrap_or(RecStatus::Pending);
1614            if let Some(f) = status_filter {
1615                if rec.status != f {
1616                    continue;
1617                }
1618            }
1619            out.push(rec);
1620        }
1621        // Review-queue order: worst first, then oldest first. Hash is only the
1622        // final tiebreak — sorting by it alone is deterministic per run but
1623        // meaningless across runs, because a grain's hash covers its timestamp,
1624        // so an identical queue comes back in a different order every time.
1625        // `dedup_key` is the last tiebreak that actually decides anything: it is
1626        // content-derived and stable across runs, whereas findings proposed in
1627        // the same sweep routinely share a `created_at_ms`. Hash trails it only
1628        // to make the ordering total.
1629        out.sort_by(|a, b| {
1630            b.severity
1631                .cmp(&a.severity)
1632                .then(a.created_at_ms.cmp(&b.created_at_ms))
1633                .then(a.dedup_key.cmp(&b.dedup_key))
1634                .then(a.hash.cmp(&b.hash))
1635        });
1636        Ok(out)
1637    }
1638
1639    /// Per-analyzer effective settings for the Setup view: the manifest facts
1640    /// merged with the file-config (override or manifest default). Read-only.
1641    pub fn analyzer_settings<S: OmsSubstrate>(
1642        &self,
1643        sub: &S,
1644    ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1645        let p = LoopPersisted::from_value(sub.load_state()?)?;
1646        Ok(self
1647            .analyzers
1648            .iter()
1649            .map(|a| {
1650                let m = a.manifest();
1651                let cfg = p.config.get(&m.id);
1652                crate::config::AnalyzerSetting {
1653                    id: m.id.clone(),
1654                    title: m.title.clone(),
1655                    description: m.description.clone(),
1656                    tier: format!("{:?}", m.tier),
1657                    trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1658                    default_on: m.default_on,
1659                    enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1660                    severity_floor: cfg
1661                        .and_then(|c| c.severity_floor)
1662                        .map(|s| s.as_str().to_string()),
1663                }
1664            })
1665            .collect())
1666    }
1667
1668    /// Update one analyzer's file-config (enable/disable, severity floor, param
1669    /// overrides, namespace scoping). Requires `Admin`. Params are validated
1670    /// against the analyzer's manifest first (unknown keys rejected, fail-closed),
1671    /// and the analyzer must exist. Returns the merged config as stored. This is
1672    /// the only write into `persisted.config` — the config layer, never a grain.
1673    pub fn set_analyzer_config<S: OmsSubstrate>(
1674        &self,
1675        sub: &mut S,
1676        analyzer_id: &str,
1677        update: crate::config::AnalyzerConfigUpdate,
1678        scopes: &ScopeSet,
1679    ) -> Result<crate::config::AnalyzerConfig> {
1680        if !scopes.has(Scope::Admin) {
1681            return Err(Error::ScopeDenied("admin".into()));
1682        }
1683        let manifest = self
1684            .analyzers
1685            .iter()
1686            .map(|a| a.manifest())
1687            .find(|m| m.id == analyzer_id)
1688            .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1689        // Validate params against the manifest BEFORE touching state.
1690        if let Some(params) = &update.params {
1691            manifest.resolve_params(params)?;
1692        }
1693        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1694        let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1695        if let Some(enabled) = update.enabled {
1696            cfg.enabled = Some(enabled);
1697        }
1698        if update.clear_floor {
1699            cfg.severity_floor = None;
1700        } else if let Some(floor) = update.severity_floor {
1701            cfg.severity_floor = Some(floor);
1702        }
1703        if let Some(params) = update.params {
1704            cfg.params = params;
1705        }
1706        if let Some(ns) = update.namespaces {
1707            cfg.namespaces = ns;
1708        }
1709        let stored = cfg.clone();
1710        sub.store_state(&p.to_value()?)?;
1711        Ok(stored)
1712    }
1713
1714    /// The measured outcome time series (the Verify gate's history) across all
1715    /// recommendations, ordered by when each checkpoint was measured.
1716    pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1717        let p = LoopPersisted::from_value(sub.load_state()?)?;
1718        let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1719        // `metric` and `rec_hash` break the tie: checkpoints measured in the
1720        // same sweep share a `measured_at_ms`, and without a tiebreak the order
1721        // falls through to the map's rec_hash ordering, which shifts every run.
1722        out.sort_by(|a, b| {
1723            a.measured_at_ms
1724                .cmp(&b.measured_at_ms)
1725                .then(a.horizon_ms.cmp(&b.horizon_ms))
1726                .then(a.metric.cmp(&b.metric))
1727                .then(a.rec_hash.cmp(&b.rec_hash))
1728        });
1729        Ok(out)
1730    }
1731
1732    /// A health snapshot — when the loop last ran, how much is un-analyzed
1733    /// since, and the queue counts. Lets a host surface "the loop may be stale"
1734    /// so a forgotten SessionEnd hook / cron doesn't silently kill it.
1735    pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
1736        let p = LoopPersisted::from_value(sub.load_state()?)?;
1737        let (grains_since_run, error_events_since_run) = count_new(sub, p.state.watermark_ms)?;
1738        let recs = self.recommendations(sub, None)?;
1739        let mut pending = 0;
1740        let mut applied = 0;
1741        for r in &recs {
1742            match r.status {
1743                RecStatus::Pending => pending += 1,
1744                RecStatus::Applied => applied += 1,
1745                _ => {}
1746            }
1747        }
1748        // Stale if it has never run, or it's been a while / a lot has piled up.
1749        let stale = match p.state.last_run_ms {
1750            None => true,
1751            Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
1752        };
1753        Ok(Health {
1754            last_run_ms: p.state.last_run_ms,
1755            grains_since_run,
1756            error_events_since_run,
1757            pending,
1758            applied,
1759            total: recs.len() as u64,
1760            stale,
1761        })
1762    }
1763
1764    /// Approval-rate metric for `origin = llm` recommendations (reflection
1765    /// design §6b) — the live field-quality signal that accrues off the audit
1766    /// chain: what fraction of the model's *surfaced* proposals a reviewer
1767    /// accepts. Complements the offline Effective-Reliability eval.
1768    pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
1769        let recs = self.recommendations(sub, None)?;
1770        let mut m = LlmMetrics::default();
1771        for r in &recs {
1772            if !matches!(r.origin, Origin::Llm { .. }) {
1773                continue;
1774            }
1775            m.proposed += 1;
1776            match r.status {
1777                RecStatus::Pending => m.pending += 1,
1778                RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
1779                RecStatus::Rejected => m.rejected += 1,
1780                RecStatus::Expired => {}
1781            }
1782        }
1783        let decided = m.approved + m.rejected;
1784        m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
1785        Ok(m)
1786    }
1787}
1788
1789/// A health snapshot for the backend's self-improvement loop.
1790#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
1791pub struct Health {
1792    #[serde(skip_serializing_if = "Option::is_none")]
1793    pub last_run_ms: Option<i64>,
1794    pub grains_since_run: u64,
1795    pub error_events_since_run: u64,
1796    pub pending: u64,
1797    pub applied: u64,
1798    pub total: u64,
1799    /// True when the loop looks stalled (never run, or ≥7d / ≥100 new grains
1800    /// since the last run) — a nudge that a trigger may be unwired.
1801    pub stale: bool,
1802}
1803
1804/// Approval-rate metric for `origin = llm` recommendations (reflection §6b).
1805#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
1806pub struct LlmMetrics {
1807    /// Total llm-origin recommendations ever stored (those that survived the
1808    /// verifier and reached the queue).
1809    pub proposed: u64,
1810    pub pending: u64,
1811    /// Approved + Applied + RolledBack (a reviewer said yes at least once).
1812    pub approved: u64,
1813    pub rejected: u64,
1814    /// approved / (approved + rejected); `None` until at least one is decided.
1815    #[serde(skip_serializing_if = "Option::is_none")]
1816    pub approval_rate: Option<f64>,
1817}
1818
1819/// Re-measure applied recommendations at each **checkpoint** past due, via the
1820/// engine's typed reads — no CAL-scalar round-trip. A recommendation
1821/// accumulates one `OutcomeResult` per horizon (measured once each), forming a
1822/// time series, so a late regression (held at 1d, regressed at 30d) is caught.
1823/// Only *regressed* checkpoints feed the outcome analyzer (→ a revert).
1824/// Unknown metric kinds are skipped, never faked.
1825fn measure_outcomes<S: OmsSubstrate>(
1826    sub: &S,
1827    p: &mut LoopPersisted,
1828    now_ms: i64,
1829) -> Result<Vec<OutcomeInput>> {
1830    // Collect all due (recommendation, horizon) checkpoints first.
1831    let mut due: Vec<(String, crate::config::AppliedRecord, i64)> = Vec::new();
1832    for (h, a) in &p.applied {
1833        if p.status_index.get(h) != Some(&RecStatus::Applied) {
1834            continue;
1835        }
1836        let Some(metric) = &a.metric else { continue };
1837        let done = p.measured.get(h).cloned().unwrap_or_default();
1838        for horizon in metric.horizons() {
1839            if now_ms - a.applied_at_ms >= horizon && !done.contains(&horizon) {
1840                due.push((h.clone(), a.clone(), horizon));
1841            }
1842        }
1843    }
1844
1845    let mut out = Vec::new();
1846    for (rec_hash, applied, horizon) in due {
1847        let metric = applied.metric.as_ref().unwrap();
1848        let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
1849            continue; // metric kind not yet re-measurable
1850        };
1851        let regressed = crate::recommendation::is_regression(
1852            metric.baseline,
1853            current,
1854            metric.higher_is_better,
1855        );
1856        p.outcomes.entry(rec_hash.clone()).or_default().push(
1857            crate::recommendation::OutcomeResult {
1858                rec_hash: rec_hash.clone(),
1859                metric: metric.metric.clone(),
1860                baseline: metric.baseline,
1861                current,
1862                verdict: if regressed { "regressed" } else { "held" }.into(),
1863                horizon_ms: horizon,
1864                measured_at_ms: now_ms,
1865            },
1866        );
1867        p.measured.entry(rec_hash.clone()).or_default().push(horizon);
1868        if regressed {
1869            out.push(OutcomeInput {
1870                rec_hash,
1871                target_ref: applied.target_ref.clone(),
1872                metric: metric.metric.clone(),
1873                baseline: metric.baseline,
1874                current,
1875                unit: metric.unit.clone(),
1876                higher_is_better: metric.higher_is_better,
1877            });
1878        }
1879    }
1880    Ok(out)
1881}
1882
1883/// Typed re-measurement for the fixed set of metric kinds the engine knows.
1884pub(crate) fn measure_metric<S: SubstrateRead>(
1885    sub: &S,
1886    metric: &crate::recommendation::MetricSnapshot,
1887    since_ms: i64,
1888) -> Result<Option<f64>> {
1889    match metric.metric.as_str() {
1890        // How many times did this tool fail again *with the same signature*
1891        // after the lesson was applied? Scoped to the signature (metric.relation)
1892        // so an unrelated later failure of the same tool is not read as a
1893        // regression of this specific lesson.
1894        "tool_error_recurrence" => {
1895            let Some(tool) = &metric.subject else { return Ok(None) };
1896            let tools = sub.grains_of_type(
1897                crate::model::grain_type::TOOL,
1898                None,
1899                ReadOpts { live_only: true, since_ms: Some(since_ms) },
1900            )?;
1901            let n = tools
1902                .iter()
1903                .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
1904                .filter(|t| {
1905                    // No stored signature (legacy metric) → fall back to the
1906                    // whole-tool count so old recommendations still measure.
1907                    metric.relation.as_deref().is_none_or(|sig| {
1908                        crate::analyzers::tool_failure::normalize_signature(
1909                            t.tool_content().unwrap_or(""),
1910                        ) == sig
1911                    })
1912                })
1913                .count();
1914            Ok(Some(n as f64))
1915        }
1916        // After a resolve-to-latest, does the subject again hold more than one
1917        // live value under the functional relation? Live-state read (no since
1918        // filter): the excess beyond one distinct object is the regression.
1919        "contradiction_recurrence" => {
1920            let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
1921                return Ok(None);
1922            };
1923            let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
1924            let distinct: BTreeSet<String> = facts
1925                .iter()
1926                .filter(|f| {
1927                    f.fact_relation()
1928                        .is_some_and(|r| normalize_ident(r) == *relation)
1929                })
1930                .filter_map(|f| f.fact_object().map(normalize_ident))
1931                .collect();
1932            Ok(Some(distinct.len().saturating_sub(1) as f64))
1933        }
1934        // `evalset:<hash>:<field>` — external correctness, measured the only
1935        // way that keeps the honesty rule ("Areev Loop improves the agent's
1936        // memory, not its outputs") intact: an evalset run is an INTERNAL,
1937        // BOUNDED, ATTRIBUTABLE measurement. The engine never runs the evalset;
1938        // it only reads summaries a host journaled with `areev eval run`.
1939        //
1940        // `since_ms` is the apply time, and it is load-bearing rather than an
1941        // optimization: a run journaled BEFORE the apply cannot be evidence of
1942        // what applying did. Comparing the baseline run against itself would
1943        // report "held" forever — a fabricated receipt, which is worse than no
1944        // receipt at all. No run since the apply → `None` → not yet
1945        // measurable, and the checkpoint stays due.
1946        m if m.starts_with("evalset:") => {
1947            let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
1948                return Ok(None);
1949            };
1950            let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
1951                return Ok(None);
1952            };
1953            // `failed`/`passed`/`total` are promoted so a metric can be written
1954            // against any evalset without the host having to add fields;
1955            // anything else is read from the summary the host did write.
1956            Ok(match field {
1957                "failed" => Some(run.failed as f64),
1958                "passed" => Some(run.passed as f64),
1959                "total" => Some(run.total() as f64),
1960                "error_rate" => match run.total() {
1961                    0 => None, // no cases ran: undefined, not zero
1962                    t => Some(run.failed as f64 / t as f64),
1963                },
1964                other => run.field(other),
1965            })
1966        }
1967        _ => Ok(None),
1968    }
1969}
1970
1971/// Live facts for one normalized (namespace?, subject) — the shared scope of
1972/// the fact-shaped recurrence metrics. `namespace: None` spans all namespaces.
1973fn scoped_live_facts<S: SubstrateRead>(
1974    sub: &S,
1975    namespace: Option<&str>,
1976    subject: &str,
1977) -> Result<Vec<GrainRecord>> {
1978    let facts = sub.grains_of_type(
1979        crate::model::grain_type::FACT,
1980        None,
1981        ReadOpts { live_only: true, since_ms: None },
1982    )?;
1983    Ok(facts
1984        .into_iter()
1985        .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
1986        .filter(|f| {
1987            f.fact_subject()
1988                .is_some_and(|s| normalize_ident(s) == subject)
1989        })
1990        .collect())
1991}
1992
1993// --- free helpers ---
1994
1995/// The §6.3 exact-equality check: a SUPERSEDE is *value-identical* when every
1996/// replacement field equals the superseded grain's value — strings after
1997/// case-fold/trim (upstream NFC is an OMS invariant), `namespace` against the
1998/// grain's own namespace, everything else exactly. This is what makes an
1999/// auto-applied consolidation provably information-preserving; a
2000/// near-duplicate (an observation body off by one token) fails it and stays
2001/// pending for human review. Fails closed: an unrecognized line shape, an
2002/// empty replacement, a missing grain, or a field the original never had all
2003/// disqualify.
2004fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
2005    let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
2006        return false;
2007    };
2008    if fields.is_empty() {
2009        return false;
2010    }
2011    let Ok(Some(grain)) = sub.grain(&target) else {
2012        return false;
2013    };
2014    // An expiry lives OUTSIDE `fields`, so the replacement (built from a fixed
2015    // field set) can never carry it — consolidating away a grain that has a
2016    // valid_to would silently drop the expiry, invisibly to the field-by-field
2017    // check below. Fail closed. (A dup that additionally carries an extra
2018    // *content* field the replacement omits is a real but subtler info-loss;
2019    // catching it soundly needs a full canonical-vs-extra comparison rather
2020    // than this replacement-scoped check, since a real fact's fields also carry
2021    // OMS metadata like `confidence` that consolidation legitimately keeps —
2022    // left as a follow-up so this narrow fix can't block valid consolidations.)
2023    if grain.valid_to_ms.is_some() {
2024        return false;
2025    }
2026    // Forward check: every replacement field equals the grain's value.
2027    fields.iter().all(|(k, v)| {
2028        if k == "namespace" {
2029            return v
2030                .as_str()
2031                .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2032        }
2033        match (v, grain.fields.get(k)) {
2034            (Value::String(a), Some(Value::String(b))) => {
2035                normalize_ident(a) == normalize_ident(b)
2036            }
2037            (a, Some(b)) => a == b,
2038            (_, None) => false,
2039        }
2040    })
2041}
2042
2043/// The DISCOVER evidence bundle's total size, and the reserved share each
2044/// source gets inside it.
2045///
2046/// Three sources feed the bundle and they answer different questions, so each
2047/// is budgeted rather than served first-come: the grains the deterministic
2048/// findings CITED (what clustering already caught), recent tool ERRORS (what
2049/// clustering could have caught but did not), and recent facts/observations —
2050/// the LLM's own lens, and the ONLY source that can carry a problem with no
2051/// error shape at all. `LENS_RESERVE` is what the last of those is guaranteed.
2052const EVIDENCE_CAP: usize = 64;
2053const CITED_SEED_CAP: usize = 24;
2054const TOOL_SEED_CAP: usize = 16;
2055/// Human-authored Observations get a small guaranteed share, taken before the
2056/// general top-up. Learning does not only come from what went wrong: a person
2057/// saying "from now on, do X" is a complete rule stated once, and no amount of
2058/// counting recovers it from a corpus that never surfaced it.
2059const NOTE_SEED_CAP: usize = 8;
2060const LENS_RESERVE: usize = 24;
2061
2062/// The confidence floor (§5.4): a verified draft below this is dropped. The
2063/// verifier's calibrated confidence is the gate, not the proposer's self-report.
2064const MIN_LLM_CONFIDENCE: f64 = 0.75;
2065
2066/// The fixed DISCOVER instruction (§5.1). The scoring rule makes "nothing to
2067/// report" a first-class, zero-penalty answer — the structural antidote to
2068/// over-generation. Kept in its own request field so it never interleaves with
2069/// (attacker-influenced) evidence text.
2070const DISCOVER_INSTRUCTIONS: &str = "You review an agent's memory for quality. \
2071Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2072findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2073stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2074recurring cost or hand-off the agent's own setup could remove). \
2075The deterministic findings already cover what the ERROR TEXT says; restating \
2076one of them earns nothing. The evidence may also contain OUTCOME records — a \
2077run's observable shape together with whether it was accepted or rejected. A \
2078problem that raised no error at all is exactly the kind the deterministic \
2079checks cannot see, so compare the rejected outcomes against the accepted \
2080ones: a feature they share and the accepted ones lack is a candidate rule. \
2081Require at least two rejected outcomes before proposing one — a single \
2082rejection is an anecdote, not a pattern. \
2083SCORING: propose a finding ONLY if you \
2084are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2085useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2086earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2087when there is nothing worth flagging. The 'approved' and 'rejected' lists, when \
2088present, show findings this reviewer recently accepted or rejected — prefer the \
2089kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2090or more evidence hashes from the bundle, name a 'target', and include your \
2091confidence 0.0-1.0. Return JSON: {\"recommendations\":[{\"summary\":\"...\",\
2092\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<hash>\"],\
2093\"confidence\":0.0,\"proposal\":{...}}]}. \
2094OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2095you are not asking to change anything. Include it ONLY when the evidence \
2096supports a specific change, choosing exactly one kind: \
2097(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2098\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
2099an action the agent itself takes on the next occasion. Either ADD an action \
2100it is failing to take ('Record the vendor name and the amount on every \
2101invoice, not just the date') or ORDER one that goes wrong ('Refund a \
2102subscription before cancelling it; refunds on cancelled subscriptions are \
2103refused'). Name the action, not a check on it: 'validate', 'verify' and \
2104'ensure ... is correct' describe a review step the agent has no way to \
2105perform, and such a rule changes nothing even once applied. \
2106(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2107entity target — a durable fact the agent keeps having to be told (an alias, a \
2108settled default, a preference). 'relation' is a short identifier (letters, \
2109digits, _ - . :), not a sentence. \
2110(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2111\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2112assembles the agent's context, when the evidence shows it retrieves the wrong \
2113things. Give the FULL new body; it replaces the old one. \
2114(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2115with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2116workflow. Only these paths are editable: 'edges.<i>.cond', \
2117'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2118holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2119(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2120— full replacement source for that tool. It is applied only after a recorded \
2121evaluation run passes, so propose one only when the evidence shows the current \
2122code is the defect. \
2123The subject of a fact, the name of a query, the plan hash and the tool name \
2124all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2125becomes a change a human reviewer may apply, so it must be fully supported by \
2126the cited evidence. Propose nothing you cannot ground in the evidence.";
2127
2128/// The fixed GROUND instruction (§5.2): verify the finding's factual PREMISES
2129/// are real (anti-fabrication), while allowing an inference. A self-improvement
2130/// finding may reason BEYOND the evidence (e.g. 'HQ=SF and country=Germany are
2131/// inconsistent'); grounding checks the premises (HQ=SF, country=Germany) are in
2132/// the evidence — the soundness of the inference is VERIFY's job, not this one.
2133const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2134fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2135the facts it relies on are actually present in the cited evidence, NOT that its \
2136conclusion is stated verbatim. Decompose the finding into the factual claims it \
2137depends on. Mark supported=true when those facts are present in the evidence \
2138(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2139on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2140different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2141
2142/// The fixed VERIFY instruction (§5.3): adversarial, abstention-biased.
2143const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2144each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2145never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2146SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2147'possible' findings with no concrete defect, and reject any claimed \
2148inconsistency or contradiction that is not backed by at least two actually \
2149conflicting facts in the cited evidence. (2) Context — does the finding \
2150correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2151finding when it names a genuine, specific problem grounded in its evidence and \
2152materially useful to a human reviewer; otherwise reject it, and default to \
2153keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2154redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2155grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2156KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2157{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2158
2159/// The fixed ENRICH instruction.
2160const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2161guidance note to help a human reviewer decide. Do not restate the finding. Return \
2162JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2163
2164/// Add a grain to the evidence bundle — deduplicated, bounded at 64, text
2165/// capped. Shared by the deterministic-citation and recent-grain seeding.
2166/// `ns_by_hash` records each bundled grain's namespace so an authored lesson
2167/// can later land in the namespace its evidence lives in (never one the
2168/// model names).
2169fn push_evidence(
2170    evidence: &mut Vec<crate::llm::EvidenceItem>,
2171    bundle: &mut BTreeSet<String>,
2172    ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2173    g: &GrainRecord,
2174) {
2175    if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2176        ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2177        evidence.push(crate::llm::EvidenceItem {
2178            hash: g.hash.clone(),
2179            grain_type: g.grain_type.clone(),
2180            text: crate::llm::cap(&grain_brief(g), 400),
2181        });
2182    }
2183}
2184
2185/// A short human-readable projection of a grain for the evidence bundle.
2186fn grain_brief(g: &GrainRecord) -> String {
2187    if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2188        return format!("{s} {r} {o}");
2189    }
2190    // Tool grains — the evidence most lesson drafts cite. Rendering them
2191    // empty starved GROUND of the very facts it exists to check: a correct
2192    // lesson would be refused as unverifiable (found live — the gate
2193    // rightly rejected a claim over evidence it could not see).
2194    if let Some(t) = g.tool_name() {
2195        let status = if g.is_error() { "error" } else { "ok" };
2196        let out = g.tool_content().unwrap_or("");
2197        return format!("tool {t} {status}: {out}");
2198    }
2199    // `object` is last but it is not optional: an Observation stores its text
2200    // there (subject + object, no relation), so it misses the fact-triple
2201    // branch above and used to fall through this list to an empty string —
2202    // every human note in a memory reached the model as a blank line. That is
2203    // the single highest-value evidence a memory holds, and it was the one
2204    // shape that rendered to nothing. Callers had started duplicating the text
2205    // into `body` to work around it; nothing should have to.
2206    for key in ["content", "body", "text", "summary", "object"] {
2207        if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2208            if !v.is_empty() {
2209                return v.to_string();
2210            }
2211        }
2212    }
2213    String::new()
2214}
2215
2216/// One imperative line, no control characters, capped: the only shape an
2217/// authored lesson may take. A literal newline could otherwise smuggle a
2218/// second CAL statement past review (belt: serde_json escapes it anyway) or
2219/// break the one-line prompt rendering hosts assume.
2220fn sanitize_lesson(s: &str) -> String {
2221    sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2222}
2223
2224/// One line, no control characters, capped. Every free-text field the model
2225/// can put into an executable proposal goes through this: a literal newline
2226/// could otherwise smuggle a second CAL statement past review (belt:
2227/// serde_json escapes it anyway), split a one-line DEFINE across the batch the
2228/// apply path iterates, or break the one-line prompt rendering hosts assume.
2229fn sanitize_line(s: &str, max: usize) -> String {
2230    let cleaned: String =
2231        s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2232    crate::llm::cap(cleaned.trim(), max)
2233}
2234
2235/// A relation is an identifier, not prose — it becomes a queryable predicate,
2236/// and whitespace or quotes in one would make the Fact unfindable by the very
2237/// recall that should surface it. `None` rejects the draft's `fact` proposal.
2238fn sanitize_relation(s: &str) -> Option<String> {
2239    let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2240    if r.is_empty()
2241        || !r
2242            .chars()
2243            .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2244    {
2245        return None;
2246    }
2247    Some(r)
2248}
2249
2250/// A model-supplied query body is placed INSIDE a `DEFINE … AS { … }` block,
2251/// so it is the one place in the vocabulary where model text becomes part of a
2252/// statement's structure rather than its content. Two belts, because the
2253/// substrate's parser strength is not something this engine gets to assume:
2254///
2255/// - **No braces.** Closing the block early is the injection shape; a saved
2256///   RECALL/ASSEMBLE body needs no braces of its own, so refusing them costs
2257///   nothing and fails closed.
2258/// - **No destructive keyword, anywhere in the body.** `cal::contains_destructive`
2259///   scans each LINE's leading keyword, which a single-line injection slips
2260///   past by construction — so scan every token here instead.
2261///
2262/// The substrate's own `validate_cal` and the saved-query read-only
2263/// verification pass still run after this; this is the layer that does not
2264/// depend on either of them being strict.
2265fn safe_definition_body(body: &str) -> bool {
2266    if body.contains('{') || body.contains('}') {
2267        return false;
2268    }
2269    !body
2270        .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2271        .any(|tok| {
2272            ["FORGET", "PURGE", "DROP", "DEFINE"]
2273                .iter()
2274                .any(|kw| tok.eq_ignore_ascii_case(kw))
2275        })
2276}
2277
2278/// The claim GROUND entails and VERIFY stress-tests. When the draft carries a
2279/// resolvable proposal the claim names exactly what an apply would do, so what
2280/// survives the gates is what gets written — never a summary standing in for
2281/// a change the gates never saw.
2282fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2283    let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2284    match resolved {
2285        Some(r) => format!("{summary} {}", r.rendered),
2286        None => summary,
2287    }
2288}
2289
2290/// A DISCOVER draft that survived structural validation, carrying the
2291/// executable form of its proposal. Resolution happens BEFORE GROUND/VERIFY,
2292/// so both gates judge exactly what an apply would do — the rule the authored
2293/// lesson already followed, generalized to the whole vocabulary. It also means
2294/// a malformed proposal costs no model call: it dies here, not at apply.
2295struct ValidatedDraft {
2296    draft: crate::llm::LlmDraft,
2297    target_ref: String,
2298    cited: Vec<String>,
2299    resolved: Option<ResolvedProposal>,
2300}
2301
2302/// The executable shape of a validated draft. `None` on a [`ValidatedDraft`]
2303/// means advisory — the model said something a human may want to see, but
2304/// nothing the engine will ever execute.
2305struct ResolvedProposal {
2306    action: ActionKind,
2307    proposal: Proposal,
2308    /// One line naming exactly what an apply would do; folded into the claim
2309    /// both gates judge and shown in the review summary.
2310    rendered: String,
2311    summary_key: &'static str,
2312    summary_args: serde_json::Map<String, Value>,
2313    rollbackable: bool,
2314    evalset_hash: Option<String>,
2315    importance: f64,
2316    /// Set for the two Fact-writing shapes. Held rather than pre-rendered so
2317    /// the grain can carry the VERIFIER's confidence, which is not known until
2318    /// after the gates have run.
2319    fact_fields: Option<serde_json::Map<String, Value>>,
2320}
2321
2322/// The Workflow fields a `plan_revision` may touch. Thresholds and limits —
2323/// never who calls what.
2324///
2325/// The exclusion is structural, not advisory: `nodes`, `edges[].src`,
2326/// `edges[].dst` and `bindings` are simply not matchable here, so a topology
2327/// change cannot be expressed by any proposal the model can write. That keeps
2328/// a plan revision reviewable as a short list of scalar deltas rather than a
2329/// re-drawn graph, which is the difference between a reviewer checking a
2330/// number and a reviewer re-deriving a plan.
2331fn plan_edit_allowed(path: &str) -> bool {
2332    let seg: Vec<&str> = path.split('.').collect();
2333    match seg.as_slice() {
2334        ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2335        ["retries", node] => !node.is_empty(),
2336        _ => false,
2337    }
2338}
2339
2340/// Read the value at an allowlisted path (absent → `Value::Null`, which is
2341/// what an edit adding a `retries` entry must declare as its `from`).
2342fn plan_get(body: &Value, path: &str) -> Value {
2343    let mut cur = body;
2344    for seg in path.split('.') {
2345        cur = match cur {
2346            Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2347                Some(v) => v,
2348                None => return Value::Null,
2349            },
2350            Value::Object(o) => match o.get(seg) {
2351                Some(v) => v,
2352                None => return Value::Null,
2353            },
2354            _ => return Value::Null,
2355        };
2356    }
2357    cur.clone()
2358}
2359
2360/// Write the value at an allowlisted path. Only creates a missing key in an
2361/// object (the `retries.<node>` case); never grows an array.
2362fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2363    let segs: Vec<&str> = path.split('.').collect();
2364    let Some((last, parents)) = segs.split_last() else {
2365        return false;
2366    };
2367    let mut cur = body;
2368    for seg in parents {
2369        cur = match cur {
2370            Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2371                Some(v) => v,
2372                None => return false,
2373            },
2374            Value::Object(o) => match o.get_mut(*seg) {
2375                Some(v) => v,
2376                None => return false,
2377            },
2378            _ => return false,
2379        };
2380    }
2381    match cur {
2382        Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2383            Some(slot) => {
2384                *slot = to;
2385                true
2386            }
2387            None => false,
2388        },
2389        Value::Object(o) => {
2390            o.insert((*last).to_string(), to);
2391            true
2392        }
2393        _ => false,
2394    }
2395}
2396
2397/// Type-check one plan edit's new value against the field it targets. Without
2398/// this a string in `max_cycles` would be dropped by the grain deserializer
2399/// and the "applied" revision would silently mean *unlimited* — a proposal
2400/// that reads as a tightening and lands as a removal.
2401fn plan_value_ok(path: &str, to: &Value) -> bool {
2402    let seg: Vec<&str> = path.split('.').collect();
2403    match seg.as_slice() {
2404        ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
2405            !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
2406        }),
2407        ["edges", _, "max_cycles"] | ["retries", _] => {
2408            to.as_u64().is_some_and(|n| n <= 1_000)
2409        }
2410        _ => false,
2411    }
2412}
2413
2414/// Resolve a DISCOVER draft's proposal into the executable form an apply would
2415/// run, or `None` for advisory. Every variant takes its SCOPE from the draft's
2416/// target and its supporting facts from the substrate — the model names the
2417/// change, never the subject it lands on, the namespace it lands in, or (for
2418/// code) the evalset that grades it.
2419fn resolve_proposal<S: OmsSubstrate>(
2420    sub: &S,
2421    d: &crate::llm::LlmDraft,
2422    target: &TargetRef,
2423    cited: &[String],
2424    ns_by_hash: &std::collections::BTreeMap<String, String>,
2425    caps: Capabilities,
2426) -> Option<ResolvedProposal> {
2427    use crate::llm::DraftProposal as P;
2428    let mut args = serde_json::Map::new();
2429    match d.parsed_proposal()? {
2430        // ---- lesson: the pre-vocabulary shape, unchanged ----
2431        P::Lesson { lesson } => {
2432            let lesson = sanitize_lesson(&lesson);
2433            if lesson.is_empty() {
2434                return None;
2435            }
2436            let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
2437            args.insert("lesson".into(), Value::from(lesson.clone()));
2438            Some(ResolvedProposal {
2439                // Same action as the deterministic lesson path — "record a
2440                // failure-derived lesson" — so dedup groups authored lessons
2441                // per target and review UIs need no new vocabulary.
2442                action: ActionKind::ClusterFailure,
2443                proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2444                rendered: format!("Proposed lesson to record: \"{lesson}\""),
2445                summary_key: "llm.lesson",
2446                summary_args: args,
2447                rollbackable: true,
2448                evalset_hash: None,
2449                importance: 0.5,
2450                fact_fields: Some(fields),
2451            })
2452        }
2453        // ---- fact: a durable fact under a model-chosen relation ----
2454        P::Fact { relation, object } => {
2455            let relation = sanitize_relation(&relation)?;
2456            let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
2457            if object.is_empty() {
2458                return None;
2459            }
2460            let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
2461            let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
2462            args.insert("relation".into(), Value::from(relation.clone()));
2463            args.insert("object".into(), Value::from(object.clone()));
2464            Some(ResolvedProposal {
2465                action: ActionKind::Record,
2466                proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2467                rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
2468                summary_key: "llm.fact",
2469                summary_args: args,
2470                rollbackable: true,
2471                evalset_hash: None,
2472                importance: 0.5,
2473                fact_fields: Some(fields),
2474            })
2475        }
2476        // ---- query_revision: how the agent assembles its own context ----
2477        P::QueryRevision { body } => {
2478            let name = target.opaque();
2479            // The name comes from the target, but it still ends up inside a
2480            // statement — a quote or control character in one would change the
2481            // statement's shape rather than its content.
2482            if name.is_empty()
2483                || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
2484            {
2485                return None;
2486            }
2487            let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
2488            if body.is_empty() || !safe_definition_body(&body) {
2489                return None;
2490            }
2491            let stmt = match target.scheme() {
2492                "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
2493                "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
2494                _ => return None,
2495            };
2496            // The substrate owns the grammar: if it will not parse, or will
2497            // not hand back an inverse, this is not something a reviewer
2498            // should be offered as applicable. A definition change ROLLBACK
2499            // could not undo must not be applied at all.
2500            sub.validate_cal(&stmt).ok()?;
2501            sub.definition_inverse(&stmt).ok().flatten()?;
2502            args.insert("name".into(), Value::from(name));
2503            args.insert("body".into(), Value::from(body.clone()));
2504            Some(ResolvedProposal {
2505                action: ActionKind::Revise,
2506                proposal: Proposal::Cal { cal: stmt },
2507                rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
2508                summary_key: "llm.query_revision",
2509                summary_args: args,
2510                rollbackable: true,
2511                evalset_hash: None,
2512                importance: 0.6,
2513                fact_fields: None,
2514            })
2515        }
2516        // ---- plan_revision: field-level edits to a Workflow grain ----
2517        P::PlanRevision { edits } => {
2518            if !caps.plans
2519                || target.scheme() != "grain"
2520                || edits.is_empty()
2521                || edits.len() > crate::llm::MAX_PLAN_EDITS
2522            {
2523                return None;
2524            }
2525            let hash = target.opaque();
2526            let g = sub.grain(hash).ok().flatten()?;
2527            if g.grain_type != "workflow" || !g.is_live() {
2528                return None;
2529            }
2530            let mut body = Value::Object(g.fields.clone());
2531            let mut deltas = Vec::new();
2532            let nodes: std::collections::BTreeSet<String> = body
2533                .get("nodes")
2534                .and_then(Value::as_array)
2535                .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
2536                .unwrap_or_default();
2537            for e in &edits {
2538                if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
2539                    return None;
2540                }
2541                // A retry count for a node that does not exist is inert, but
2542                // applying it still mints a new plan hash — and every trigger
2543                // pointing at the old one must then be walked forward. A
2544                // no-op is not worth that.
2545                if let Some(node) = e.path.strip_prefix("retries.") {
2546                    if !nodes.contains(node) {
2547                        return None;
2548                    }
2549                }
2550                // `from` is the staleness check: a proposal authored against
2551                // an older plan does not silently apply to a newer one.
2552                if plan_get(&body, &e.path) != e.from {
2553                    return None;
2554                }
2555                // A no-op edit is not a revision; it would apply, mint a new
2556                // plan hash, and orphan every trigger pointing at the old one
2557                // for nothing.
2558                if e.from == e.to {
2559                    return None;
2560                }
2561                if !plan_set(&mut body, &e.path, e.to.clone()) {
2562                    return None;
2563                }
2564                deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
2565            }
2566            // The substrate owns the plan grammar (unique + reachable nodes,
2567            // conditions parse, every cycle bounded). An edit that would make
2568            // the plan unrunnable never reaches a reviewer as applicable.
2569            sub.validate_plan(&body).ok()?;
2570            let Value::Object(fields) = body else {
2571                return None;
2572            };
2573            let stmt = cal::supersede(hash, "workflow", &fields);
2574            // Same rule as the definition rewrite: a statement the substrate
2575            // will not accept is not something to offer a reviewer as
2576            // applicable. `validate_plan` checked the GRAPH; this checks the
2577            // statement that carries it.
2578            sub.validate_cal(&stmt).ok()?;
2579            args.insert("plan".into(), Value::from(hash));
2580            args.insert("edits".into(), Value::from(deltas.join("; ")));
2581            Some(ResolvedProposal {
2582                action: ActionKind::Revise,
2583                proposal: Proposal::Cal { cal: stmt },
2584                rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
2585                summary_key: "llm.plan_revision",
2586                summary_args: args,
2587                rollbackable: true,
2588                evalset_hash: None,
2589                importance: 0.7,
2590                fact_fields: None,
2591            })
2592        }
2593        // ---- code_revision: §7.4, gated by the tool's own evalset ----
2594        P::CodeRevision { source } => {
2595            if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
2596                return None;
2597            }
2598            if source.chars().count() > crate::llm::MAX_CODE_LEN {
2599                return None;
2600            }
2601            // Rule E1's pin, resolved from the substrate. A proposer that
2602            // could name its own grader is not gated, and a tool that
2603            // declares no evalset has no gate to pass — advisory either way.
2604            let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
2605            let mut data = serde_json::Map::new();
2606            data.insert("tool".into(), Value::from(target.opaque()));
2607            data.insert("source".into(), Value::from(source.clone()));
2608            args.insert("tool".into(), Value::from(target.opaque()));
2609            args.insert("bytes".into(), Value::from(source.len() as u64));
2610            Some(ResolvedProposal {
2611                action: ActionKind::CodeRevision,
2612                proposal: Proposal::Data { data },
2613                rendered: format!(
2614                    "Proposed new source for tool {} ({} bytes), gated by evalset {}",
2615                    target.opaque(),
2616                    source.len(),
2617                    evalset
2618                ),
2619                summary_key: "llm.code_revision",
2620                summary_args: args,
2621                rollbackable: true,
2622                evalset_hash: Some(evalset),
2623                importance: 0.8,
2624                fact_fields: None,
2625            })
2626        }
2627    }
2628}
2629
2630/// Stamp a validated DISCOVER draft as an `origin = llm` recommendation.
2631/// Default shape: an advisory `Flag` carrying `Proposal::Data` (no executable
2632/// mutation). A draft whose proposal RESOLVED (see [`resolve_proposal`])
2633/// instead stamps as that executable change, reviewable with the exact line an
2634/// apply would run. Either way `Origin::Llm` plus the no-manifest analyzer id
2635/// leave it structurally ineligible for auto-apply — and independently, no
2636/// class this vocabulary can reach except `memory` is auto-appliable at all
2637/// (`Policy::grants_auto_apply`). The only path into the agent is a human
2638/// review with a BECAUSE followed by an explicit apply.
2639fn stamp_llm(
2640    model: &str,
2641    d: &crate::llm::LlmDraft,
2642    target_ref: String,
2643    cited: Vec<String>,
2644    resolved: Option<ResolvedProposal>,
2645    confidence: f64,
2646    now_ms: i64,
2647) -> Recommendation {
2648    let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2649    let guidance = if d.guidance.trim().is_empty() {
2650        None
2651    } else {
2652        Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
2653    };
2654    let (action, proposal, summary, rollbackable, importance, evalset_hash) = match resolved {
2655        Some(mut r) => {
2656            // The grain records the VERIFIER's calibrated confidence — the
2657            // independent signal — never the proposer's self-report.
2658            if let Some(mut fields) = r.fact_fields.take() {
2659                fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
2660                r.proposal = Proposal::Cal { cal: cal::add("fact", &fields) };
2661            }
2662            let mut args = r.summary_args;
2663            args.insert("text".into(), Value::from(summary_text));
2664            (
2665                r.action,
2666                r.proposal,
2667                Summary::new(r.summary_key, args),
2668                r.rollbackable,
2669                r.importance,
2670                r.evalset_hash,
2671            )
2672        }
2673        None => {
2674            let mut args = serde_json::Map::new();
2675            args.insert("text".into(), Value::from(summary_text));
2676            let mut data = serde_json::Map::new();
2677            data.insert("source".into(), Value::from("llm"));
2678            (
2679                ActionKind::Flag,
2680                Proposal::Data { data },
2681                Summary::new("llm.discover", args),
2682                false,
2683                0.3,
2684                None,
2685            )
2686        }
2687    };
2688    Recommendation {
2689        hash: String::new(),
2690        analyzer: "loop.llm/1".to_string(),
2691        params_snapshot: serde_json::Map::new(),
2692        origin: Origin::Llm { model: model.to_string() },
2693        target_ref: target_ref.clone(),
2694        action_kind: action,
2695        dedup_key: dedup_key("llm", &target_ref, action),
2696        summary,
2697        severity: Severity::Low,
2698        proposal,
2699        destructive: false,
2700        rollbackable,
2701        evidence: cited,
2702        evidence_query: None,
2703        metric: None,
2704        // The verifier's calibrated confidence — not a hardcoded default.
2705        confidence: confidence.clamp(0.0, 1.0),
2706        importance,
2707        created_at_ms: now_ms,
2708        guidance,
2709        evalset_hash,
2710        status: RecStatus::Pending,
2711    }
2712}
2713
2714/// The Fact an approved `lesson` or `fact` proposal writes: subject from the
2715/// entity target, the model's relation (`"lesson"` for the lesson shape —
2716/// prescriptive prose, distinct from the deterministic `fails_with` signature
2717/// facts), the sanitized object, and the DOMINANT namespace of the cited
2718/// evidence (max count, ties to the lexicographically smallest — the
2719/// tool_failure rule), never a namespace the model names. `None` when the
2720/// target gives no subject.
2721fn derived_fact_fields(
2722    target: &TargetRef,
2723    relation: &str,
2724    object: &str,
2725    cited: &[String],
2726    ns_by_hash: &std::collections::BTreeMap<String, String>,
2727) -> Option<serde_json::Map<String, Value>> {
2728    if target.scheme() != "entity" {
2729        return None;
2730    }
2731    let subject = target
2732        .opaque()
2733        .rsplit_once('/')
2734        .map(|(_, s)| s)
2735        .unwrap_or(target.opaque());
2736    if subject.is_empty() {
2737        return None;
2738    }
2739    let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
2740    for h in cited {
2741        if let Some(ns) = ns_by_hash.get(h) {
2742            if !ns.is_empty() {
2743                *ns_counts.entry(ns.as_str()).or_default() += 1;
2744            }
2745        }
2746    }
2747    let lesson_ns = ns_counts
2748        .iter()
2749        .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
2750        .map(|(ns, _)| ns.to_string());
2751    let mut fields = serde_json::Map::new();
2752    fields.insert("subject".into(), Value::from(subject));
2753    fields.insert("relation".into(), Value::from(relation));
2754    fields.insert("object".into(), Value::from(object));
2755    // `confidence` is stamped by the caller from the VERIFIER's calibrated
2756    // score, not the proposer's self-report — so it is deliberately absent
2757    // here, where only the proposer has spoken.
2758    if let Some(ns) = lesson_ns {
2759        fields.insert("namespace".into(), Value::from(ns));
2760    }
2761    Some(fields)
2762}
2763
2764/// The action kinds that apply ONLY through the evalset-run gating edge
2765/// (§7.4 for tool code; the tuning seam's adapter promotion inherits the
2766/// same rule). One predicate so the gate, the rollbackable stamp, and the
2767/// promotion write can never disagree on membership.
2768fn requires_gating(kind: ActionKind) -> bool {
2769    matches!(
2770        kind,
2771        ActionKind::CodeRevision | ActionKind::AdapterRevision
2772    )
2773}
2774
2775fn stamp(
2776    m: &AnalyzerManifest,
2777    params: &crate::manifest::Params,
2778    d: crate::recommendation::RecDraft,
2779    now_ms: i64,
2780) -> Result<Recommendation> {
2781    let target = TargetRef::parse(&d.target_ref)?;
2782    // Rule E1 at the door (§7.4): an unpinned code revision, a pinned
2783    // non-code target, or a code action on a non-tool target never becomes
2784    // a recommendation at all.
2785    crate::recommendation::validate_code_rules(
2786        d.action_kind,
2787        target.target_class(),
2788        d.evalset_hash.as_deref(),
2789    )?;
2790    let dedup = dedup_key(m.family(), &d.target_ref, d.action_kind);
2791    let destructive = match &d.proposal {
2792        Proposal::Cal { cal } => cal::contains_destructive(cal),
2793        _ => false,
2794    };
2795    let rollbackable = match &d.proposal {
2796        Proposal::Cal { .. } => !destructive,
2797        Proposal::Edit { .. } => false,
2798        // A code or adapter revision applies by WRITING the promotion
2799        // grain; retracting it is the exact inverse — rollbackable by
2800        // construction.
2801        Proposal::Data { .. } => requires_gating(d.action_kind),
2802    };
2803    let mut evidence = d.evidence;
2804    evidence.truncate(MAX_EVIDENCE);
2805    // Provenance follows the analyzer's trust class: a subprocess
2806    // (`--analyzer-cmd`) finding is stamped `Command` — structurally
2807    // auto-apply-ineligible and badged [external] on the recall surface —
2808    // not `Builtin`.
2809    let origin = match m.trust_class {
2810        crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
2811        _ => Origin::Builtin,
2812    };
2813    Ok(Recommendation {
2814        hash: String::new(),
2815        analyzer: m.id.clone(),
2816        params_snapshot: params.snapshot(),
2817        origin,
2818        target_ref: target.as_string(),
2819        action_kind: d.action_kind,
2820        dedup_key: dedup,
2821        summary: d.summary,
2822        severity: d.severity,
2823        proposal: d.proposal,
2824        destructive,
2825        rollbackable,
2826        evidence,
2827        evidence_query: d.evidence_query,
2828        metric: d.metric,
2829        confidence: d.confidence,
2830        importance: d.importance,
2831        created_at_ms: now_ms,
2832        guidance: None,
2833        evalset_hash: d.evalset_hash,
2834        status: RecStatus::Pending,
2835    })
2836}
2837
2838fn validate_because(because: &str) -> Result<String> {
2839    let trimmed = because.trim();
2840    if trimmed.is_empty() {
2841        return Err(Error::InvalidProposal(
2842            "a BECAUSE reason is required".into(),
2843        ));
2844    }
2845    if trimmed.chars().count() > MAX_BECAUSE {
2846        return Err(Error::InvalidProposal(format!(
2847            "BECAUSE exceeds {MAX_BECAUSE} chars"
2848        )));
2849    }
2850    Ok(trimmed.to_string())
2851}
2852
2853fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
2854    for req in &m.requires {
2855        match req {
2856            Capability::Forks if !caps.forks => return Some("forks"),
2857            Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
2858            Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
2859            _ => {}
2860        }
2861    }
2862    None
2863}
2864
2865fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
2866    p.config.get(analyzer_id).and_then(|c| c.severity_floor)
2867}
2868
2869fn gate(
2870    opts: &RunOptions,
2871    p: &LoopPersisted,
2872    new_grains: u64,
2873    new_errors: u64,
2874    now_ms: i64,
2875) -> Option<SkipReason> {
2876    let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
2877    if !any {
2878        return None;
2879    }
2880    let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
2881    let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
2882    let stale_ok = opts
2883        .if_stale_ms
2884        .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
2885    if min_new_ok || min_err_ok || stale_ok {
2886        return None;
2887    }
2888    // Choose the honest reason: staleness-only gate → not_stale, else min_new.
2889    if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
2890        Some(SkipReason::NotStale)
2891    } else {
2892        Some(SkipReason::MinNewNotMet)
2893    }
2894}
2895
2896fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<(u64, u64)> {
2897    let opts = ReadOpts {
2898        live_only: false,
2899        since_ms: watermark.map(|w| w + 1),
2900    };
2901    let mut new_grains = 0u64;
2902    let mut new_errors = 0u64;
2903    for t in [
2904        crate::model::grain_type::FACT,
2905        crate::model::grain_type::EVENT,
2906        crate::model::grain_type::TOOL,
2907        crate::model::grain_type::OBSERVATION,
2908    ] {
2909        let g = sub.grains_of_type(t, None, opts)?;
2910        new_grains += g.len() as u64;
2911        // The error gate (--min-new-errors) watches captured tool failures.
2912        if t == crate::model::grain_type::TOOL {
2913            new_errors += g.iter().filter(|e| e.is_error()).count() as u64;
2914        }
2915    }
2916    Ok((new_grains, new_errors))
2917}
2918
2919fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
2920    let grains = sub.grains_of_type(
2921        crate::model::grain_type::RECOMMENDATION,
2922        Some(LOOP_NS),
2923        ReadOpts {
2924            live_only: false,
2925            since_ms: None,
2926        },
2927    )?;
2928    let mut set = BTreeSet::new();
2929    for g in grains {
2930        let status = p
2931            .status_index
2932            .get(&g.hash)
2933            .copied()
2934            .unwrap_or(RecStatus::Pending);
2935        // Pending/approved (still open) and applied (already handled)
2936        // recommendations suppress re-proposal of the same finding. Rejected
2937        // is handled by cooldowns; rolled_back/expired may legitimately
2938        // re-propose (the situation returned).
2939        if matches!(
2940            status,
2941            RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
2942        ) {
2943            if let Some(key) = g.str_field("dedup_key") {
2944                set.insert(key.to_string());
2945            }
2946        }
2947    }
2948    Ok(set)
2949}
2950
2951fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
2952    let g = sub
2953        .grain(rec_hash)?
2954        .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
2955    Recommendation::from_fields(rec_hash, &g.fields)
2956}
2957
2958/// Is this line a definition rewrite — a statement that changes a saved
2959/// `qry:`/`tpl:` registry row rather than writing a grain?
2960///
2961/// A keyword test, not a parse: the engine deliberately contains a CAL
2962/// *writer*, never a parser (parsing is the substrate's job). Both spellings
2963/// are matched case-insensitively, and `DROP` is intentionally absent — the
2964/// loop may propose defining a query, never removing one.
2965pub(crate) fn is_definition_statement(line: &str) -> bool {
2966    let up = line.trim_start().to_ascii_uppercase();
2967    up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
2968}
2969
2970/// The refusal an advisory `Edit` earns. The engine has no executable edit
2971/// primitive; the change belongs in the host.
2972const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
2973     Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
2974     approve it to acknowledge it and let it expire.";
2975
2976/// The refusal an advisory `Data` finding earns — every `Data` shape except
2977/// `outcome_review`'s revert, which carries `revert_of`.
2978/// Shared by [`Engine::preflight_apply`] and the apply gate so a fused
2979/// approve-and-apply caller is refused BEFORE the approval lands, not after.
2980const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
2981     (evalset hash + run id + stats) — use apply_gated";
2982
2983const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
2984     guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
2985     acknowledge it and let it expire.";
2986
2987/// Whether [`Engine::apply`] can execute this proposal at all.
2988///
2989/// One source of truth, shared by [`Engine::preflight_apply`] and
2990/// [`Engine::apply`] so the two can never disagree.
2991///
2992/// Deliberately **not** consulted on approve. An advisory finding is a `Flag`:
2993/// approving it means "yes, this is real", which is the whole workflow for the
2994/// LLM path and the telemetry analyzers. What must not happen is a caller being
2995/// walked into an approval and *then* refused — which is exactly what the fused
2996/// approve-and-apply path in the bindings did, leaving the recommendation in
2997/// `approved`, whose only exits are `applied` and `expired`. Preflight asks
2998/// first, so that path now refuses before it commits anything.
2999pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
3000    match proposal {
3001        Proposal::Cal { .. } => Ok(()),
3002        Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
3003        // Two executable Data shapes: `outcome_review`'s `revert_of` (names
3004        // an earlier applied recommendation to roll back), and a gated
3005        // revision (code or adapter), which executes by writing its
3006        // promotion grain — the gating-run requirement itself is checked at
3007        // apply, not here.
3008        Proposal::Data { data } => {
3009            if requires_gating(action_kind)
3010                || data.get("revert_of").and_then(Value::as_str).is_some()
3011            {
3012                Ok(())
3013            } else {
3014                Err(Error::InvalidProposal(ADVISORY_DATA.into()))
3015            }
3016        }
3017    }
3018}
3019
3020#[cfg(test)]
3021mod definition_body_tests {
3022    use super::safe_definition_body;
3023
3024    #[test]
3025    fn ordinary_bodies_pass() {
3026        assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
3027        assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
3028    }
3029
3030    #[test]
3031    fn a_body_cannot_close_its_own_block_or_carry_destruction() {
3032        // The injection shape: close the DEFINE block, append a statement.
3033        // Newlines are already collapsed by `sanitize_line`, so the payload
3034        // arrives as ONE line — which is exactly what a line-leading-keyword
3035        // destructive scan cannot see.
3036        assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
3037        assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
3038        // …and the keyword alone is refused even without the braces.
3039        assert!(!safe_definition_body("RECALL facts FORGET abc"));
3040        assert!(!safe_definition_body("recall facts purge older than 1d"));
3041        assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
3042        // A nested DEFINE would redefine something the target does not name.
3043        assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
3044        // Substring matches are not keywords — this must still pass.
3045        assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
3046    }
3047}
3048
3049#[cfg(test)]
3050mod plan_edit_tests {
3051    use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
3052    use serde_json::json;
3053
3054    fn plan() -> serde_json::Value {
3055        json!({
3056            "nodes": ["fetch", "review", "post"],
3057            "edges": [
3058                {"src": "fetch", "dst": "review"},
3059                {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
3060            ],
3061            "bindings": {"fetch": "sha256:tool1"},
3062            "retries": {"fetch": 1}
3063        })
3064    }
3065
3066    #[test]
3067    fn the_allowlist_admits_thresholds_and_refuses_topology() {
3068        assert!(plan_edit_allowed("edges.1.cond"));
3069        assert!(plan_edit_allowed("edges.1.max_cycles"));
3070        assert!(plan_edit_allowed("retries.fetch"));
3071        // Topology is not expressible — the structural half of the guarantee
3072        // that a plan revision stays reviewable as scalar deltas.
3073        for path in [
3074            "nodes",
3075            "nodes.0",
3076            "edges.0.src",
3077            "edges.0.dst",
3078            "edges",
3079            "bindings.fetch",
3080            "edges.x.cond",
3081            "",
3082        ] {
3083            assert!(!plan_edit_allowed(path), "{path} must not be editable");
3084        }
3085    }
3086
3087    #[test]
3088    fn values_are_type_checked_against_the_field() {
3089        // A string in max_cycles would be DROPPED by the grain deserializer,
3090        // so an "applied" tightening would silently mean unlimited.
3091        assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
3092        assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
3093        assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
3094        assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
3095        assert!(plan_value_ok("retries.fetch", &json!(3)));
3096        assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
3097        assert!(!plan_value_ok("edges.1.cond", &json!("  ")));
3098        assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
3099        assert!(!plan_value_ok("edges.0.src", &json!("other")));
3100    }
3101
3102    #[test]
3103    fn get_reads_through_arrays_and_objects_and_absence_is_null() {
3104        let p = plan();
3105        assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
3106        assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
3107        // An edit that ADDS a retry declares `from: null` — so absence has to
3108        // read as Null rather than as an error.
3109        assert_eq!(plan_get(&p, "retries.review"), json!(null));
3110        assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
3111        assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
3112    }
3113
3114    #[test]
3115    fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
3116        let mut p = plan();
3117        assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
3118        assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
3119        assert!(plan_set(&mut p, "retries.review", json!(2)));
3120        assert_eq!(plan_get(&p, "retries.review"), json!(2));
3121        assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
3122        assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
3123    }
3124}
3125
3126#[cfg(test)]
3127mod definition_proposal_tests {
3128    use super::is_definition_statement;
3129
3130    #[test]
3131    fn definition_statements_are_recognized_in_both_spellings() {
3132        assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
3133        assert!(is_definition_statement("  define template foo AS { x }"));
3134        assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
3135        // Ordinary proposals are untouched.
3136        assert!(!is_definition_statement("ADD fact {}"));
3137        assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
3138        assert!(!is_definition_statement("FORGET abc"));
3139        // `DROP` is never proposable, so it is deliberately NOT a definition
3140        // statement here — a proposal containing one still fails validation
3141        // rather than being handed an inverse.
3142        assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
3143    }
3144}