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