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, Checkpoint};
19use crate::substrate::{Capabilities, OmsSubstrate, ReadOpts, SubstrateRead};
20use serde::{Deserialize, Serialize};
21use serde_json::{Map, Value};
22use std::collections::{BTreeMap, BTreeSet};
23
24/// The namespace the loop's own grains (recommendations, audit) live in.
25pub const LOOP_NS: &str = "areev-loop";
26
27/// Host-granted authority, per connection. `admin` implies all.
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum Scope {
30    Read,
31    Write,
32    Review,
33    Apply,
34    Admin,
35}
36
37/// A set of granted scopes.
38#[derive(Debug, Clone, Default)]
39pub struct ScopeSet(Vec<Scope>);
40
41impl ScopeSet {
42    pub fn of(scopes: &[Scope]) -> Self {
43        ScopeSet(scopes.to_vec())
44    }
45    /// The local root of trust: whoever can run against the file holds all
46    /// scopes (the CLI/embedded posture).
47    pub fn all() -> Self {
48        ScopeSet(vec![Scope::Admin])
49    }
50    pub fn has(&self, s: Scope) -> bool {
51        self.0.contains(&Scope::Admin) || self.0.contains(&s)
52    }
53}
54
55/// A review decision.
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
57pub enum Decision {
58    Approve,
59    Reject,
60}
61
62/// Gating and scoping options for a run.
63#[derive(Debug, Clone, Default)]
64pub struct RunOptions {
65    pub min_new: Option<u64>,
66    pub min_new_errors: Option<u64>,
67    pub if_stale_ms: Option<i64>,
68    /// Optional global namespace filter (empty = all).
69    pub namespaces: Vec<String>,
70    /// Re-analyze the WHOLE memory this pass, not just grains since the last-run
71    /// watermark (`areev loop reflect`). Dedup/cooldowns still suppress anything
72    /// already queued, and the watermark still advances at the end — so a sweep
73    /// is safe to run any time and later runs stay incremental. Mainly widens the
74    /// watermark-sensitive inputs (tool-failure window, the non-parasitic LLM
75    /// evidence bundle) to the full history.
76    pub full_sweep: bool,
77    /// The principal that invoked this run. Recorded as co-creator on every
78    /// non-`Builtin` (LLM / external-command) recommendation the run stores,
79    /// so the review gate's self-approval block also fires for whoever
80    /// triggered the model that authored the finding — an LLM draft is
81    /// authored *via* its trigger, unlike a deterministic finding, which is
82    /// computed. `None` (a headless/scheduled run) records no co-creator.
83    pub triggering_actor: Option<String>,
84}
85
86/// Whether a run executed or was skipped.
87#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
88#[serde(rename_all = "lowercase")]
89pub enum RunOutcome {
90    Ran,
91    Skipped,
92}
93
94/// Why a run was a no-op. `LockHeld` is produced by the host adapter (a
95/// concurrent writer), surfaced here for a single contract.
96#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
97#[serde(rename_all = "snake_case")]
98pub enum SkipReason {
99    MinNewNotMet,
100    NotStale,
101    LockHeld,
102    /// The host policy's `cadence` block set a threshold and none was met.
103    CadenceNotDue,
104}
105
106/// One analyzer that did not contribute drafts, with why.
107#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
108pub struct AnalyzerSkip {
109    pub id: String,
110    pub reason: String,
111}
112
113/// The run-outcome contract (proposal §13): one shape across CLI/API/MCP/bindings.
114#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
115pub struct RunResult {
116    pub outcome: RunOutcome,
117    #[serde(skip_serializing_if = "Option::is_none")]
118    pub skip_reason: Option<SkipReason>,
119    pub new_grains: u64,
120    pub new_error_events: u64,
121    pub proposed: u64,
122    pub deduped: u64,
123    pub stored: u64,
124    /// Of the stored recommendations, how many were auto-applied by policy.
125    #[serde(default)]
126    pub auto_applied: u64,
127    #[serde(default)]
128    pub analyzers_run: Vec<String>,
129    #[serde(default)]
130    pub analyzers_skipped: Vec<AnalyzerSkip>,
131    /// Where the LLM's contribution went, stage by stage. `None` when no
132    /// backend is attached.
133    #[serde(default, skip_serializing_if = "Option::is_none")]
134    pub llm_funnel: Option<LlmFunnel>,
135}
136
137/// The DISCOVER pipeline's attrition, counted.
138///
139/// "The model contributed nothing" has at least five distinct causes, and
140/// they call for opposite responses: an empty bundle is a capture problem, a
141/// model that abstained may need better evidence or a better prompt, drafts
142/// dying at GROUND suggest fabrication, drafts dying at VERIFY suggest they
143/// were vague, and drafts dying at the floor were merely unconfident. Without
144/// this they are indistinguishable from the outside — every one of them
145/// renders as an empty ledger and reads like a clean null. That ambiguity
146/// cost a six-cell measurement run before it was noticed.
147#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
148pub struct LlmFunnel {
149    /// Grains offered to the model as evidence. Zero means nothing to reflect on.
150    pub evidence: u64,
151    /// Drafts the model returned.
152    pub proposed: u64,
153    /// Survived the cite-check and target-class filter.
154    pub cited: u64,
155    /// Dropped for citing a hash the bundle does not hold. Its own counter
156    /// because the fix is specific: models copy long hex badly, and that is
157    /// a different problem from one aiming at a surface it may not touch.
158    pub dropped_uncited: u64,
159    /// Dropped for targeting a class the vocabulary may not reach — a prompt,
160    /// host config, or its own grader.
161    pub dropped_target: u64,
162    /// Survived GROUND — their premises were found in the cited evidence.
163    pub grounded: u64,
164    /// Verdicts GROUND actually returned. `grounded = 0` with verdicts > 0 is
165    /// the gate refusing every draft; `grounded = 0` with verdicts = 0 is a
166    /// grader that answered with nothing usable, which is a backend problem
167    /// wearing a gate's clothes.
168    pub ground_verdicts: u64,
169    /// The GROUND call itself failed — no response at all. The engine
170    /// fail-softs here by design, so without this the run looks like a model
171    /// that had nothing to say.
172    pub ground_call_failed: bool,
173    /// Survived VERIFY's adversarial pass.
174    pub kept: u64,
175    /// Cleared the confidence floor and reached the queue.
176    pub stored: u64,
177    /// Kept, but demoted to advisory for citing fewer distinct evidence
178    /// grains than `Policy::min_evidence`. Omitted when zero, so a run under
179    /// the default policy reads exactly as before.
180    #[serde(default, skip_serializing_if = "is_zero")]
181    pub advisory_thin_evidence: u64,
182    /// Kept, but dropped before the queue for restating a live lesson on the
183    /// same entity in other words (`Policy::near_duplicate = suppress`). Its
184    /// own counter beside `dropped_uncited` and `dropped_target`, so "the
185    /// model contributed nothing" keeps its distinct causes. Omitted when
186    /// zero.
187    #[serde(default, skip_serializing_if = "is_zero")]
188    pub dropped_near_duplicate: u64,
189}
190
191fn is_zero(n: &u64) -> bool {
192    *n == 0
193}
194
195impl RunResult {
196    fn skipped(reason: SkipReason, new_grains: u64, new_error_events: u64) -> Self {
197        RunResult {
198            outcome: RunOutcome::Skipped,
199            skip_reason: Some(reason),
200            new_grains,
201            new_error_events,
202            proposed: 0,
203            deduped: 0,
204            stored: 0,
205            auto_applied: 0,
206            llm_funnel: None,
207            analyzers_run: vec![],
208            analyzers_skipped: vec![],
209        }
210    }
211
212    pub fn ran(&self) -> bool {
213        self.outcome == RunOutcome::Ran
214    }
215}
216
217/// The engine holds the registered analyzers, the host policy, and an optional
218/// LLM enrichment backend (§9).
219pub struct Engine {
220    analyzers: Vec<Box<dyn Analyzer>>,
221    policy: crate::policy::Policy,
222    /// Optional LLM backend. `None` → the DISCOVER/ENRICH stages are the
223    /// identity, so the pipeline is byte-for-byte the deterministic path.
224    llm: Option<Box<dyn crate::llm::LlmBackend>>,
225    /// Optional separate backend for the GROUND stage (§5.2, §11). `None` →
226    /// grounding rides `llm`. Lets a team point entailment at a cheaper or
227    /// specialized model (or take the generative model out of grounding
228    /// entirely) without changing the proposer/verifier.
229    ground_llm: Option<Box<dyn crate::llm::LlmBackend>>,
230}
231
232pub(crate) struct AnalysisPass {
233    pub(crate) survivors: Vec<Recommendation>,
234    proposed: u64,
235    deduped: u64,
236    analyzers_run: Vec<String>,
237    pub(crate) analyzers_skipped: Vec<AnalyzerSkip>,
238    llm_funnel: Option<LlmFunnel>,
239}
240
241impl Engine {
242    /// An engine with the default built-ins and a default (fully closed)
243    /// policy — nothing auto-applies, no LLM.
244    pub fn with_builtins() -> Self {
245        Engine {
246            analyzers: crate::analyzer::builtin_analyzers(),
247            policy: crate::policy::Policy::default(),
248            llm: None,
249            ground_llm: None,
250        }
251    }
252
253    /// An engine with no analyzers (register your own).
254    pub fn empty() -> Self {
255        Engine {
256            analyzers: vec![],
257            policy: crate::policy::Policy::default(),
258            llm: None,
259            ground_llm: None,
260        }
261    }
262
263    /// Install a host policy (the only place auto-apply is granted).
264    pub fn with_policy(mut self, policy: crate::policy::Policy) -> Self {
265        self.policy = policy;
266        self
267    }
268
269    /// Attach an optional LLM enrichment backend (§9). Only ever *adds* cited
270    /// draft recommendations (stamped `origin = llm`, never auto-applied) and
271    /// whitelisted guidance notes — it can never gate or rewrite deterministic
272    /// output.
273    pub fn with_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
274        self.llm = Some(backend);
275        self
276    }
277
278    /// Attach a separate backend for the GROUND stage (§5.2). Without this,
279    /// grounding uses the `with_llm` backend. Independent of the proposer so an
280    /// operator can run entailment on a cheaper/specialized model.
281    pub fn with_ground_llm(mut self, backend: Box<dyn crate::llm::LlmBackend>) -> Self {
282        self.ground_llm = Some(backend);
283        self
284    }
285
286    pub fn policy(&self) -> &crate::policy::Policy {
287        &self.policy
288    }
289
290    /// Whether an LLM backend is attached (replay reports it as not replayed).
291    pub fn has_llm(&self) -> bool {
292        self.llm.is_some()
293    }
294
295    /// Register an additional analyzer (the linked-Rust seam).
296    pub fn register(&mut self, analyzer: Box<dyn Analyzer>) {
297        self.analyzers.push(analyzer);
298    }
299
300    pub fn analyzers(&self) -> &[Box<dyn Analyzer>] {
301        &self.analyzers
302    }
303
304    /// Run the exact production analysis/validation path without Phase 0
305    /// measurement or Phase 3 persistence. The immutable substrate borrow is
306    /// the replay safety boundary: recommendations, audit grains, state,
307    /// cooldowns, outcomes, and the op-log cannot be changed here.
308    ///
309    /// `overrides` is keyed by full analyzer id and overlays the file's stored
310    /// parameter map. Unknown keys fail closed through `resolve_params` and
311    /// surface as an analyzer skip, exactly as in a production run.
312    pub fn analyze_only<S: OmsSubstrate>(
313        &self,
314        sub: &S,
315        opts: &RunOptions,
316        overrides: &BTreeMap<String, Map<String, Value>>,
317        now_ms: i64,
318    ) -> Result<Vec<Recommendation>> {
319        let persisted = LoopPersisted::from_value(sub.load_state()?)?;
320        let analysis_watermark = if opts.full_sweep {
321            None
322        } else {
323            persisted.state.watermark_ms
324        };
325        Ok(self
326            .analysis_pass(
327                sub,
328                &persisted,
329                opts,
330                overrides,
331                analysis_watermark,
332                now_ms,
333                &[],
334            )?
335            .survivors)
336    }
337
338    /// Run one analysis pass. Idempotent under `dedup_key`; the watermark is
339    /// advanced at the end, so a crashed run simply re-runs.
340    pub fn run<S: OmsSubstrate>(
341        &self,
342        sub: &mut S,
343        opts: &RunOptions,
344        now_ms: i64,
345    ) -> Result<RunResult> {
346        let mut persisted = LoopPersisted::from_value(sub.load_state()?)?;
347        let watermark = persisted.state.watermark_ms;
348        // A full sweep analyzes the whole memory (watermark ignored for the
349        // analysis inputs), while gating, `new` counts, and the end-of-run
350        // watermark advance still use the real watermark. Dedup/cooldowns keep
351        // it from re-proposing what is already queued.
352        let analysis_watermark = if opts.full_sweep { None } else { watermark };
353
354        let new = count_new(sub, watermark)?;
355        let (new_grains, new_error_events) = (new.grains, new.error_events);
356        if let Some(reason) = gate(opts, &persisted, new_grains, new_error_events, now_ms) {
357            return Ok(RunResult::skipped(reason, new_grains, new_error_events));
358        }
359        // The policy's cadence applies when the caller set no gate of its own
360        // (host CLI flags > policy file) and is not asking for a sweep — a
361        // sweep is a command, not a tick.
362        let flags_set = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
363        if !flags_set && !opts.full_sweep {
364            if let Some(reason) = cadence_gate(&self.policy.cadence, &persisted, new, now_ms) {
365                return Ok(RunResult::skipped(reason, new_grains, new_error_events));
366            }
367        }
368
369        // Phase 0: re-measure applied recommendations due for review (the
370        // Verify gate). Records a measured outcome per due recommendation —
371        // and, under policy, asks the gate's second question: does each
372        // applied recommendation's PREMISE still stand?
373        let mut outcome_inputs = measure_outcomes(sub, &mut persisted, &self.policy, now_ms)?;
374        if self.policy.premise_drift {
375            outcome_inputs.extend(detect_premise_drift(sub, &mut persisted, now_ms)?);
376        }
377
378        let AnalysisPass {
379            survivors,
380            proposed,
381            deduped,
382            analyzers_run,
383            analyzers_skipped,
384            llm_funnel,
385        } = self.analysis_pass(
386            &*sub,
387            &persisted,
388            opts,
389            &BTreeMap::new(),
390            analysis_watermark,
391            now_ms,
392            &outcome_inputs,
393        )?;
394
395        // Phase 3 (needs &mut): store survivors + propose audit, then
396        // auto-apply the ones the host policy grants (all gates in §6.3).
397        let mut stored = 0u64;
398        let mut auto_applied = 0u64;
399        for mut rec in survivors {
400            let spec = rec.to_grain_spec(LOOP_NS)?;
401            let hash = sub.put_grain(&spec)?;
402            rec.hash = hash.clone();
403            let actor = format!("engine:{}", rec.analyzer);
404            let audit = AuditRecord {
405                rec_hash: hash.clone(),
406                from: None,
407                to: RecStatus::Pending,
408                actor: actor.clone(),
409                observer_type: ObserverType::System,
410                because: "analyzer proposed".into(),
411                previous_audit_hash: None,
412                gating: None,
413                at_ms: now_ms,
414            };
415            let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
416            persisted
417                .status_index
418                .insert(hash.clone(), RecStatus::Pending);
419            persisted.creators.insert(hash.clone(), actor);
420            // An LLM or external-command finding exists because someone ran
421            // it — record that principal too, so review can refuse the
422            // trigger approving their own model's output. Builtin analyzers
423            // stay engine-only: deterministic output has no human author.
424            if !matches!(rec.origin, Origin::Builtin) {
425                if let Some(trigger) = &opts.triggering_actor {
426                    persisted.co_creators.insert(hash.clone(), trigger.clone());
427                }
428            }
429            persisted.audit_heads.insert(hash.clone(), audit_hash);
430            stored += 1;
431
432            if self.can_auto_apply(&*sub, &rec) {
433                self.auto_apply(sub, &mut persisted, &rec, now_ms)?;
434                auto_applied += 1;
435            }
436        }
437
438        persisted.state.last_run_ms = Some(now_ms);
439        persisted.state.watermark_ms = Some(now_ms);
440        sub.store_state(&persisted.to_value()?)?;
441
442        Ok(RunResult {
443            outcome: RunOutcome::Ran,
444            skip_reason: None,
445            new_grains,
446            new_error_events,
447            proposed,
448            deduped,
449            stored,
450            auto_applied,
451            analyzers_run,
452            analyzers_skipped,
453            llm_funnel,
454        })
455    }
456
457    /// Shared Phase 1–2 implementation for production and replay. Cooldowns
458    /// and the live recommendation queue are honored in both modes; only the
459    /// production caller proceeds into the mutating Phase 3 below.
460    #[allow(clippy::too_many_arguments)]
461    #[allow(clippy::too_many_arguments)]
462    fn analysis_pass<S: OmsSubstrate>(
463        &self,
464        sub: &S,
465        persisted: &LoopPersisted,
466        opts: &RunOptions,
467        external_overrides: &BTreeMap<String, Map<String, Value>>,
468        analysis_watermark: Option<i64>,
469        now_ms: i64,
470        outcome_inputs: &[OutcomeInput],
471    ) -> Result<AnalysisPass> {
472        let existing = existing_dedup_keys(sub, persisted)?;
473        self.analysis_pass_inner(
474            sub,
475            persisted,
476            &self.policy,
477            opts,
478            external_overrides,
479            analysis_watermark,
480            now_ms,
481            outcome_inputs,
482            &existing,
483            None,
484        )
485    }
486
487    /// The pass itself. `existing` is the set of dedup keys already open
488    /// (production computes it from the stored queue; replay carries its own
489    /// simulated queue). `replay` is `Some(reason)` when the pass is a
490    /// rehearsal: the LLM stage and every analyzer that is not a pure
491    /// function of the grains (external commands, telemetry rollups) are
492    /// skipped with that reason rather than run against the present.
493    #[allow(clippy::too_many_arguments)]
494    pub(crate) fn analysis_pass_inner<S: OmsSubstrate>(
495        &self,
496        sub: &S,
497        persisted: &LoopPersisted,
498        policy: &crate::policy::Policy,
499        opts: &RunOptions,
500        external_overrides: &BTreeMap<String, Map<String, Value>>,
501        analysis_watermark: Option<i64>,
502        now_ms: i64,
503        outcome_inputs: &[OutcomeInput],
504        existing: &BTreeSet<String>,
505        replay: Option<&str>,
506    ) -> Result<AnalysisPass> {
507        let mut analyzers_run = Vec::new();
508        let mut analyzers_skipped = Vec::new();
509        let mut candidates: Vec<Recommendation> = Vec::new();
510        let caps = sub.capabilities();
511        let verdicts = latest_verdicts(persisted);
512
513        for analyzer in &self.analyzers {
514            let m = analyzer.manifest();
515            if let Some(why) = replay {
516                let out_of_process = m.trust_class == crate::manifest::TrustClass::Command;
517                let telemetry_fed = m.requires.contains(&crate::manifest::Capability::Telemetry);
518                if out_of_process || telemetry_fed {
519                    analyzers_skipped.push(AnalyzerSkip {
520                        id: m.id.clone(),
521                        reason: format!(
522                            "not replayed: {}",
523                            if out_of_process { why } else { "telemetry rollups are not time-indexed" }
524                        ),
525                    });
526                    continue;
527                }
528            }
529            let cfg = persisted.config.get(&m.id);
530            let enabled = cfg.and_then(|c| c.enabled).unwrap_or(m.default_on);
531            if !enabled {
532                analyzers_skipped.push(AnalyzerSkip {
533                    id: m.id.clone(),
534                    reason: "disabled".into(),
535                });
536                continue;
537            }
538            if policy.denies(m.family()) {
539                analyzers_skipped.push(AnalyzerSkip {
540                    id: m.id.clone(),
541                    reason: "denied by host policy".into(),
542                });
543                continue;
544            }
545            if let Some(missing) = missing_capability(m, caps) {
546                analyzers_skipped.push(AnalyzerSkip {
547                    id: m.id.clone(),
548                    reason: format!("missing capability: {missing}"),
549                });
550                continue;
551            }
552            let mut param_overrides = cfg.map(|c| c.params.clone()).unwrap_or_default();
553            if let Some(extra) = external_overrides.get(&m.id) {
554                for (key, value) in extra {
555                    param_overrides.insert(key.clone(), value.clone());
556                }
557            }
558            let params = match m.resolve_params(&param_overrides) {
559                Ok(p) => p,
560                Err(e) => {
561                    analyzers_skipped.push(AnalyzerSkip {
562                        id: m.id.clone(),
563                        reason: e.to_string(),
564                    });
565                    continue;
566                }
567            };
568            let ns_owned = cfg.map(|c| c.namespaces.clone()).unwrap_or_default();
569            let ns_slice: &[String] = if ns_owned.is_empty() {
570                &opts.namespaces
571            } else {
572                &ns_owned
573            };
574            let reader: &dyn SubstrateRead = sub;
575            let ctx = AnalyzeCtx::new(
576                reader,
577                &params,
578                ns_slice,
579                analysis_watermark,
580                now_ms,
581                outcome_inputs,
582                &verdicts,
583            );
584            match analyzer.analyze(&ctx) {
585                Ok(drafts) => {
586                    analyzers_run.push(m.id.clone());
587                    for draft in drafts {
588                        match stamp(m, &params, draft, now_ms) {
589                            Ok(rec) => candidates.push(rec),
590                            Err(e) => analyzers_skipped.push(AnalyzerSkip {
591                                id: m.id.clone(),
592                                reason: e.to_string(),
593                            }),
594                        }
595                    }
596                }
597                Err(e) => analyzers_skipped.push(AnalyzerSkip {
598                    id: m.id.clone(),
599                    reason: e.to_string(),
600                }),
601            }
602        }
603
604        let mut funnel = LlmFunnel::default();
605        if self.llm.is_some() && replay.is_none() {
606            candidates.extend(self.discover(
607                sub,
608                &candidates,
609                analysis_watermark,
610                &opts.namespaces,
611                now_ms,
612                &mut funnel,
613            ));
614        }
615
616        let proposed = candidates.len() as u64;
617        let mut seen = BTreeSet::new();
618        let mut survivors = Vec::new();
619        for candidate in candidates {
620            let family = crate::manifest::analyzer_family(&candidate.analyzer);
621            let floor = [
622                severity_floor_for(persisted, &candidate.analyzer),
623                policy.severity_floor(family),
624            ]
625            .into_iter()
626            .flatten()
627            .max();
628            if floor.is_some_and(|floor| candidate.severity < floor) {
629                continue;
630            }
631            if !seen.insert(candidate.dedup_key.clone()) {
632                continue;
633            }
634            if existing.contains(&candidate.dedup_key) {
635                continue;
636            }
637            if persisted
638                .cooldowns
639                .get(&candidate.dedup_key)
640                .is_some_and(|until| now_ms < *until)
641            {
642                continue;
643            }
644            survivors.push(candidate);
645        }
646        let deduped = proposed - survivors.len() as u64;
647        if self.llm.is_some() && replay.is_none() {
648            self.enrich(&mut survivors);
649        }
650        Ok(AnalysisPass {
651            survivors,
652            proposed,
653            deduped,
654            analyzers_run,
655            analyzers_skipped,
656            llm_funnel: self.llm.is_some().then_some(funnel),
657        })
658    }
659
660    /// DISCOVER (§9): ask the LLM for additional draft recommendations, given
661    /// the deterministic findings as *context* and a bounded, provenance-tagged
662    /// evidence bundle. Every returned draft must cite evidence present in the
663    /// bundle and target a memory/query surface; it is stamped `origin = llm`
664    /// (so it can never auto-apply) and enters the ordinary dedup/store path. A
665    /// failed or garbled response yields no drafts — never a failed run.
666    fn discover<S: OmsSubstrate>(
667        &self,
668        sub: &S,
669        candidates: &[Recommendation],
670        watermark: Option<i64>,
671        namespaces: &[String],
672        now_ms: i64,
673        funnel: &mut LlmFunnel,
674    ) -> Vec<Recommendation> {
675        let Some(llm) = &self.llm else {
676            return Vec::new();
677        };
678        let findings: Vec<crate::llm::FindingBrief> = candidates
679            .iter()
680            .take(32)
681            .map(|c| crate::llm::FindingBrief {
682                analyzer: c.analyzer.clone(),
683                summary: c.summary.render(),
684                target: c.target_ref.clone(),
685                severity: c.severity.as_str().to_string(),
686            })
687            .collect();
688        // Evidence bundle. Seeded first from the grains the deterministic
689        // findings cite, THEN — the non-parasitic step (§11) — topped up with
690        // RECENT grains (created since the last run) so the LLM gets its own
691        // lens and can find issues in grains no analyzer flagged. Without this
692        // the LLM could only elaborate near what determinism already caught.
693        let attribution = self.policy.evidence_attribution;
694        let mut evidence: Vec<crate::llm::EvidenceItem> = Vec::new();
695        let mut bundle: BTreeSet<String> = BTreeSet::new();
696        let mut ns_by_hash: std::collections::BTreeMap<String, String> = Default::default();
697        'cited: for c in candidates {
698            for h in &c.evidence {
699                // Every source gets a RESERVED share of the bundle, because a
700                // cap that only one source respects is not a budget. A single
701                // tool_failure finding may cite up to MAX_EVIDENCE (64)
702                // grains — the whole bundle — so without this the "give the
703                // LLM its own lens" seeding below could be starved to nothing
704                // by the very determinism it is supposed to look past. Found
705                // live: with the deterministic findings citing enough, the
706                // model never saw a single non-cited grain.
707                if evidence.len() >= CITED_SEED_CAP {
708                    break 'cited;
709                }
710                if !bundle.contains(h) {
711                    if let Ok(Some(g)) = sub.grain(h) {
712                        push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
713                    }
714                }
715            }
716        }
717        let scan_ns: Vec<Option<&str>> = if namespaces.is_empty() {
718            vec![None]
719        } else {
720            namespaces.iter().map(|n| Some(n.as_str())).collect()
721        };
722        let opts = ReadOpts { live_only: true, since_ms: watermark };
723        // Tool grains carry the raw experience of a tool-using agent, and an
724        // ERROR is the part a reflection pass can act on. Seeded after the
725        // cited grains and inside its own reserved share, so a busy desk
726        // cannot crowd out the facts and observations below.
727        //
728        // Without this the LLM saw tool failures only through the
729        // deterministic findings that happened to cite them: it could
730        // elaborate on what clustering already caught, but could never find
731        // a failure clustering missed — the one thing it is here for. The
732        // top-up below called itself "non-parasitic" while omitting the very
733        // grain type the flagship analyzer reads.
734        //
735        // Two passes over the same reserved share: failures first, then — only
736        // when the host lets the proposer author skills — the calls that
737        // SUCCEEDED. A skill's evidence is a trajectory that worked, and a
738        // model shown only what broke can never propose one (measured: with a
739        // successful procedure in the memory and no failure, the bundle was
740        // empty and DISCOVER was never called). Errors keep first claim on the
741        // share, so a busy desk's successes cannot bury the failure signal the
742        // lesson path exists for.
743        let mut tool_seeded = 0usize;
744        let seed_successes = self.policy.skills.enabled || self.policy.plans.enabled;
745        'tools: for want_error in [true, false] {
746            if !want_error && !seed_successes {
747                break;
748            }
749            for ns in &scan_ns {
750                if let Ok(recent) = sub.grains_of_type(crate::model::grain_type::TOOL, *ns, opts) {
751                    for g in recent {
752                        if tool_seeded >= TOOL_SEED_CAP
753                            || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE
754                        {
755                            break 'tools;
756                        }
757                        if g.is_error() != want_error {
758                            continue;
759                        }
760                        let before = evidence.len();
761                        push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
762                        if evidence.len() > before {
763                            tool_seeded += 1;
764                        }
765                    }
766                }
767            }
768        }
769        // Harness evidence, ahead of the general passes and read by EXPLICIT
770        // namespace — the one path that sees past the `agent:` exclusion, and
771        // only for the kinds named above. Early and reserved for the same
772        // reason human notes are: a fold summary is written once per long run
773        // and would lose every time to routine volume.
774        //
775        // A missing namespace or no read grant is nothing to reflect on, not an
776        // error — the same posture `run_outcome` takes over the same rows.
777        if let Ok(rows) =
778            sub.grains_of_type(crate::model::grain_type::OBSERVATION, Some(crate::eval::HARNESS_NS), opts)
779        {
780            let mut seeded = 0usize;
781            for g in rows {
782                if seeded >= HARNESS_SEED_CAP || evidence.len() >= EVIDENCE_CAP - LENS_RESERVE {
783                    break;
784                }
785                let kind = g.str_field("observation_kind").unwrap_or_default();
786                if !HARNESS_EVIDENCE_KINDS.contains(&kind) {
787                    continue;
788                }
789                let before = evidence.len();
790                push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
791                if evidence.len() > before {
792                    seeded += 1;
793                }
794            }
795        }
796        // Observations BEFORE facts, with their own small reserve.
797        //
798        // An Observation is where a human's own words land — a supervisor's
799        // note, a reviewer's correction, an instruction for next time. That
800        // signal is stated ONCE by nature, so recency and frequency seeding
801        // structurally bury it: one note loses to three hundred routine
802        // records every time, and the rarest evidence is usually the most
803        // valuable. Exhausting facts first (as this did) meant a desk with
804        // any volume showed the model no human input at all.
805        'notes: for ns in &scan_ns {
806            if let Ok(recent) =
807                sub.grains_of_type(crate::model::grain_type::OBSERVATION, *ns, opts)
808            {
809                for g in recent {
810                    if evidence.len() >= CITED_SEED_CAP + TOOL_SEED_CAP + NOTE_SEED_CAP {
811                        break 'notes;
812                    }
813                    push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
814                }
815            }
816        }
817        'seed: for gt in [
818            crate::model::grain_type::FACT,
819            crate::model::grain_type::OBSERVATION,
820        ] {
821            for ns in &scan_ns {
822                if let Ok(recent) = sub.grains_of_type(gt, *ns, opts) {
823                    for g in recent {
824                        if evidence.len() >= EVIDENCE_CAP {
825                            break 'seed;
826                        }
827                        push_evidence(&mut evidence, &mut bundle, &mut ns_by_hash, &g, attribution);
828                    }
829                }
830            }
831        }
832        funnel.evidence = evidence.len() as u64;
833        if evidence.is_empty() {
834            return Vec::new(); // nothing to reflect on
835        }
836        // PROPOSE (§5.1): the abstention-legitimate objective — "nothing to
837        // report" is a first-class, zero-penalty answer. The operator-taste
838        // history (recent approve/reject decisions on llm findings) is passed so
839        // the model learns what this reviewer accepts.
840        let (approved, rejected) = self.llm_history(sub);
841        let base = match self.policy.discover_objective {
842            crate::policy::DiscoverObjective::ReviewQueue => DISCOVER_INSTRUCTIONS,
843            crate::policy::DiscoverObjective::Learner => DISCOVER_LEARNER_INSTRUCTIONS,
844        };
845        // The skill kind is offered only when the host allows it, so the
846        // vocabulary the model sees is exactly the vocabulary that can apply.
847        let mut instructions = base.to_string();
848        if self.policy.skills.enabled {
849            instructions.push_str(&skill_instructions(self.policy.skills.min_steps));
850        }
851        if self.policy.plans.enabled {
852            instructions.push_str(&plan_instructions(self.policy.plans.min_nodes));
853        }
854        // Offered only when a pile is on the table: the paragraph would
855        // otherwise invite the model to invent one.
856        if findings.iter().any(|f| f.analyzer.starts_with("loop.lesson_pile/")) {
857            instructions.push_str(CONSOLIDATION_INSTRUCTIONS);
858        }
859        let request = crate::llm::LlmRequest {
860            loop_proto: 1,
861            op: "discover",
862            instructions: &instructions,
863            findings: findings.clone(),
864            evidence: evidence.clone(),
865            rejected,
866            approved,
867        };
868        let Ok(body) = serde_json::to_string(&request) else {
869            return Vec::new();
870        };
871        let raw = match llm.complete(&body) {
872            Ok(r) => r,
873            Err(_) => return Vec::new(), // fail-soft
874        };
875        // Cheap structural validation (cite-check + target class); collect the
876        // survivors for the verifier. Storing the normalized target string
877        // avoids a TargetRef clone through the pipeline.
878        let caps = sub.capabilities();
879        let mut validated: Vec<ValidatedDraft> = Vec::new();
880        let drafts: Vec<_> = crate::llm::parse_discover(&raw)
881            .recommendations
882            .into_iter()
883            .take(crate::llm::MAX_LLM_DRAFTS)
884            .collect();
885        funnel.proposed = drafts.len() as u64;
886        let id_to_hash: std::collections::BTreeMap<&str, &str> = evidence
887            .iter()
888            .map(|e| (e.id.as_str(), e.hash.as_str()))
889            .collect();
890        for d in drafts {
891            let mut cited: Vec<String> = Vec::new();
892            for c in &d.evidence {
893                if let Some(h) = resolve_citation(c, &bundle, &id_to_hash) {
894                    if !cited.contains(&h) {
895                        cited.push(h);
896                    }
897                }
898            }
899            if cited.is_empty() {
900                funnel.dropped_uncited += 1;
901                continue; // uncited → drop (no fabrication)
902            }
903            let Ok(target) = TargetRef::parse(&d.target) else {
904                funnel.dropped_target += 1;
905                continue;
906            };
907            let tc = target.target_class();
908            // The classes the proposal vocabulary can reach: memory
909            // (entity/grain), query (query/template) and code (tool). The
910            // prompt (`doc:`), `host:`, `evalset:` and `model:` classes stay
911            // closed to the model — it may not rewrite the agent's prompt, its
912            // host config, or the gate that grades its own code.
913            if !matches!(tc, "memory" | "query" | "code") {
914                funnel.dropped_target += 1;
915                continue;
916            }
917            // Resolve BEFORE the gates: what GROUND entails and VERIFY
918            // stress-tests is exactly what an apply would do. A draft citing
919            // fewer distinct grains than the host's evidence floor is kept
920            // as a finding but offered as nothing a reviewer could apply — a
921            // single instance may be worth a person's attention; it is not,
922            // under that policy, a rule.
923            let thin = cited.len() < self.policy.min_evidence as usize;
924            if thin {
925                funnel.advisory_thin_evidence += 1;
926            }
927            let resolved = if thin {
928                None
929            } else {
930                resolve_proposal(sub, &d, &target, &cited, &ns_by_hash, caps, &self.policy)
931            };
932            // A `tool:` target has exactly one legal shape (Rule E1: a code
933            // target REQUIRES action_kind code_revision), so an unresolved one
934            // could not even be stamped advisory — drop it here rather than
935            // spend two model calls on something that fails validation after.
936            if tc == "code" && resolved.is_none() {
937                funnel.dropped_target += 1;
938                continue;
939            }
940            validated.push(ValidatedDraft {
941                draft: d,
942                target_ref: target.as_string(),
943                cited,
944                resolved,
945            });
946        }
947        funnel.cited = validated.len() as u64;
948        if validated.is_empty() {
949            return Vec::new();
950        }
951        // GROUND → VERIFY → ROUTE (§5.2–5.4): only drafts that survive an
952        // independent grounding entailment check *and* an adversarial
953        // verification pass (each a separate call — proposer ≠ scorer) reach the
954        // queue, stamped with the verifier's calibrated confidence.
955        // GROUND may run on a separate backend (§11); VERIFY always uses the
956        // main llm (the proposer≠scorer independence is on VERIFY, not GROUND).
957        let ground = self.ground_llm.as_deref().unwrap_or(&**llm);
958        let outcome_metric = self.outcome_metric_template(sub);
959        self.verify_drafts(sub, &**llm, ground, validated, &evidence, outcome_metric, now_ms, funnel)
960    }
961
962    /// The metric an applicable LLM-authored proposal will be re-measured by,
963    /// when the host policy names an evalset (`Policy::outcome_evalset`).
964    /// The baseline is the newest run journaled so far — the state of the
965    /// world BEFORE the proposal, which is what "did applying it help" has
966    /// to be read against. No run journaled yet → no metric: the lesson is
967    /// honestly unmeasured rather than scored against a number nobody
968    /// recorded.
969    fn outcome_metric_template<S: OmsSubstrate>(
970        &self,
971        sub: &S,
972    ) -> Option<crate::recommendation::MetricSnapshot> {
973        let e = self.policy.outcome_evalset.as_ref()?;
974        let run = crate::eval::newest_eval_run(sub, &e.hash, None).ok().flatten()?;
975        let baseline = crate::eval::run_value(&run, &e.field)?;
976        // The schedule in the host's unit. `review_after_ms`/`horizons_ms` keep
977        // the time view for anything that only understands time; when the
978        // host counts runs or grains, `checkpoints` is the schedule.
979        let schedule = e.schedule();
980        let ms_only: Vec<i64> = schedule.iter().filter_map(Checkpoint::as_ms).collect();
981        let all_ms = ms_only.len() == schedule.len();
982        let horizons = if ms_only.is_empty() { vec![86_400_000] } else { ms_only };
983        Some(crate::recommendation::MetricSnapshot {
984            metric: format!("evalset:{}:{}", e.hash, e.field),
985            baseline,
986            unit: e.field.clone(),
987            n: run.total(),
988            window: "per-run".into(),
989            subject: None,
990            namespace: None,
991            relation: None,
992            query: format!(
993                "RECALL facts WHERE subject = \"evalset:{}\" AND relation = \"mg:eval_run\"",
994                e.hash
995            ),
996            review_after_ms: horizons[0],
997            horizons_ms: if all_ms { horizons } else { Vec::new() },
998            checkpoints: if all_ms { Vec::new() } else { schedule },
999            higher_is_better: e.higher_is_better,
1000        })
1001    }
1002
1003    /// GROUND → VERIFY → ROUTE (§5.2–5.4). Two independent model calls, batched
1004    /// over the drafts: a grounding-entailment gate ("does the cited evidence
1005    /// support the claim?"), then an adversarial keep/kill with a calibrated
1006    /// confidence. A draft reaches the queue only if it is grounded **and** kept
1007    /// **and** clears the confidence floor. Any failed call drops the whole LLM
1008    /// contribution for the run (safe default), never the run.
1009    #[allow(clippy::too_many_arguments)]
1010    #[allow(clippy::too_many_arguments)]
1011    fn verify_drafts<S: SubstrateRead>(
1012        &self,
1013        sub: &S,
1014        llm: &dyn crate::llm::LlmBackend,
1015        ground: &dyn crate::llm::LlmBackend,
1016        validated: Vec<ValidatedDraft>,
1017        evidence: &[crate::llm::EvidenceItem],
1018        outcome_metric: Option<crate::recommendation::MetricSnapshot>,
1019        now_ms: i64,
1020        funnel: &mut LlmFunnel,
1021    ) -> Vec<Recommendation> {
1022        use crate::llm::*;
1023        let ev_by_hash: std::collections::BTreeMap<&str, &EvidenceItem> =
1024            evidence.iter().map(|e| (e.hash.as_str(), e)).collect();
1025        let ev_for = |cited: &[String]| -> Vec<EvidenceItem> {
1026            cited
1027                .iter()
1028                .filter_map(|h| ev_by_hash.get(h.as_str()).map(|e| (*e).clone()))
1029                .collect()
1030        };
1031
1032        // GROUND (§5.2): decompose-then-entail per draft, batched into one call.
1033        // The claim includes any authored lesson — the gate must judge exactly
1034        // what an apply would record, not only the finding's summary.
1035        let claims: Vec<GroundItem> = validated
1036            .iter()
1037            .enumerate()
1038            .map(|(i, v)| GroundItem {
1039                id: i,
1040                claim: claim_text(&v.draft, v.resolved.as_ref()),
1041                evidence: ev_for(&v.cited),
1042            })
1043            .collect();
1044        let ground_req = GroundRequest {
1045            loop_proto: 1,
1046            op: "ground",
1047            instructions: GROUND_INSTRUCTIONS,
1048            claims,
1049        };
1050        // A grounding pass that REFUSED every draft and one that never
1051        // answered are the same number of survivors and opposite problems:
1052        // the first is the gate doing its job, the second is a backend having
1053        // a bad minute while the engine fail-softs. Count the verdicts
1054        // actually returned so the two are distinguishable afterwards.
1055        let grounded: std::collections::BTreeSet<usize> = match serde_json::to_string(&ground_req)
1056            .ok()
1057            .and_then(|b| ground.complete(&b).ok())
1058        {
1059            Some(raw) => {
1060                let parsed = parse_ground(&raw);
1061                funnel.ground_verdicts = parsed.results.len() as u64;
1062                parsed
1063                    .results
1064                    .into_iter()
1065                    .filter(|r| r.supported)
1066                    .map(|r| r.id)
1067                    .collect()
1068            }
1069            None => {
1070                funnel.ground_call_failed = true;
1071                return Vec::new();
1072            }
1073        };
1074        funnel.grounded = grounded.len() as u64;
1075        if grounded.is_empty() {
1076            return Vec::new();
1077        }
1078
1079        // VERIFY (§5.3): adversarial keep/kill over the grounded drafts, a
1080        // separate call from the proposer. Soundness + abstention only — NOT
1081        // novelty. Novelty is steered at DISCOVER and settled by human review;
1082        // asking a weak verifier to judge it just makes it hallucinate "already
1083        // known" and kill genuine findings (§11).
1084        let items: Vec<VerifyItem> = validated
1085            .iter()
1086            .enumerate()
1087            .filter(|(i, _)| grounded.contains(i))
1088            .map(|(i, v)| VerifyItem {
1089                id: i,
1090                // Same rule as GROUND: the adversarial pass sees the change.
1091                summary: claim_text(&v.draft, v.resolved.as_ref()),
1092                target: v.target_ref.clone(),
1093                evidence: ev_for(&v.cited),
1094            })
1095            .collect();
1096        let verify_req = VerifyRequest {
1097            loop_proto: 1,
1098            op: "verify",
1099            instructions: VERIFY_INSTRUCTIONS,
1100            findings: items,
1101        };
1102        let verdicts: std::collections::BTreeMap<usize, f64> =
1103            match serde_json::to_string(&verify_req).ok().and_then(|b| llm.complete(&b).ok()) {
1104                Some(raw) => parse_verify(&raw)
1105                    .results
1106                    .into_iter()
1107                    .filter(|r| r.keep)
1108                    .map(|r| (r.id, r.confidence.clamp(0.0, 1.0)))
1109                    .collect(),
1110                None => return Vec::new(),
1111            };
1112
1113        funnel.kept = verdicts.len() as u64;
1114        // ROUTE (§5.4): grounded ∧ kept ∧ verifier-confidence ≥ floor. The
1115        // verifier's confidence (the independent signal) is what we trust and
1116        // stamp — not the proposer's self-report.
1117        let mut out = Vec::new();
1118        for (i, v) in validated.into_iter().enumerate() {
1119            if let Some(&conf) = verdicts.get(&i) {
1120                if conf >= MIN_LLM_CONFIDENCE {
1121                    // A lesson that restates, in other words, a live lesson
1122                    // on the same entity. `authored_dedup_key` collapses the
1123                    // same TEXT; meaning is measured here — cosine over the
1124                    // substrate's embedder when it has one, token-set Jaccard
1125                    // as the T0 floor — and the policy says whether the
1126                    // reviewer sees it marked or never sees it. A
1127                    // consolidation is exempt: superseding the pile is its
1128                    // whole point.
1129                    let near = match v.resolved.as_ref() {
1130                        Some(r) if r.action != ActionKind::Consolidate => r
1131                            .fact_fields
1132                            .as_ref()
1133                            .filter(|f| f.get("relation").and_then(Value::as_str) == Some("lesson"))
1134                            .map(|f| {
1135                                near_duplicates_of(
1136                                    sub,
1137                                    f.get("subject").and_then(Value::as_str).unwrap_or(""),
1138                                    f.get("namespace").and_then(Value::as_str),
1139                                    f.get("object").and_then(Value::as_str).unwrap_or(""),
1140                                )
1141                            })
1142                            .unwrap_or_default(),
1143                        _ => Vec::new(),
1144                    };
1145                    if !near.is_empty()
1146                        && self.policy.near_duplicate == crate::policy::NearDuplicateMode::Suppress
1147                    {
1148                        funnel.dropped_near_duplicate += 1;
1149                        continue;
1150                    }
1151                    let mut rec = stamp_llm(
1152                        llm.model(),
1153                        &v.draft,
1154                        v.target_ref,
1155                        v.cited,
1156                        v.resolved,
1157                        conf,
1158                        now_ms,
1159                    );
1160                    if let Some(best) = near.first() {
1161                        // The summary says so, and names the closest rule.
1162                        rec.summary.args.insert("near_count".into(), Value::from(near.len() as u64));
1163                        rec.summary.args.insert("near_score".into(), Value::from(best.score));
1164                        rec.summary.args.insert("near_method".into(), Value::from(best.method.clone()));
1165                        rec.summary.args.insert(
1166                            "near_hash".into(),
1167                            Value::from(best.hash.chars().take(12).collect::<String>()),
1168                        );
1169                        rec.summary.template_id = "llm.lesson_near_duplicate".into();
1170                        rec.near_duplicate_of = near;
1171                    }
1172                    // Only a proposal an apply can execute (and roll back)
1173                    // is measured: an advisory flag changes nothing, so
1174                    // there is nothing to hold or regress.
1175                    if rec.rollbackable {
1176                        rec.metric = outcome_metric.clone();
1177                    }
1178                    out.push(rec);
1179                }
1180            }
1181        }
1182        funnel.stored = out.len() as u64;
1183        out
1184    }
1185
1186    /// Recent operator decisions on `origin = llm` findings — approved (incl.
1187    /// applied) and rejected summaries, most-recent first and bounded — so
1188    /// DISCOVER can learn what this reviewer accepts (§9). Best-effort: a read
1189    /// failure yields empty history, never an error.
1190    fn llm_history<S: OmsSubstrate>(&self, sub: &S) -> (Vec<String>, Vec<String>) {
1191        const MAX: usize = 20;
1192        let Ok(mut recs) = self.recommendations(sub, None) else {
1193            return (Vec::new(), Vec::new());
1194        };
1195        recs.retain(|r| matches!(r.origin, Origin::Llm { .. }));
1196        recs.sort_by_key(|r| std::cmp::Reverse(r.created_at_ms));
1197        let mut approved = Vec::new();
1198        let mut rejected = Vec::new();
1199        for r in &recs {
1200            match r.status {
1201                RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack
1202                    if approved.len() < MAX =>
1203                {
1204                    approved.push(r.summary.render());
1205                }
1206                RecStatus::Rejected if rejected.len() < MAX => rejected.push(r.summary.render()),
1207                _ => {}
1208            }
1209        }
1210        (approved, rejected)
1211    }
1212
1213    /// ENRICH (§9): ask the LLM to add a short guidance note to the surviving
1214    /// deterministic recommendations. Whitelist-only — only `guidance` is
1215    /// merged (capped), and only onto recs that don't already have one; the
1216    /// engine-templated summary is never touched. Fail-soft.
1217    fn enrich(&self, survivors: &mut [Recommendation]) {
1218        let Some(llm) = &self.llm else {
1219            return;
1220        };
1221        if survivors.is_empty() {
1222            return;
1223        }
1224        let findings: Vec<crate::llm::FindingBrief> = survivors
1225            .iter()
1226            .map(|r| crate::llm::FindingBrief {
1227                analyzer: r.analyzer.clone(),
1228                summary: r.summary.render(),
1229                target: r.target_ref.clone(),
1230                severity: r.severity.as_str().to_string(),
1231            })
1232            .collect();
1233        let request = crate::llm::LlmRequest {
1234            loop_proto: 1,
1235            op: "enrich",
1236            instructions: ENRICH_INSTRUCTIONS,
1237            findings,
1238            evidence: Vec::new(),
1239            rejected: Vec::new(),
1240            approved: Vec::new(),
1241        };
1242        let Ok(body) = serde_json::to_string(&request) else {
1243            return;
1244        };
1245        let raw = match llm.complete(&body) {
1246            Ok(r) => r,
1247            Err(_) => return,
1248        };
1249        for note in crate::llm::parse_enrich(&raw).notes {
1250            if note.guidance.trim().is_empty() {
1251                continue;
1252            }
1253            if let Some(r) = survivors
1254                .iter_mut()
1255                .find(|r| r.target_ref == note.target && r.guidance.is_none())
1256            {
1257                r.guidance = Some(crate::llm::cap(&note.guidance, crate::llm::MAX_GUIDANCE_LEN));
1258            }
1259        }
1260    }
1261
1262    /// Evaluate the auto-apply gate (§6.3) — ALL preconditions must hold:
1263    /// host opt-in + policy grant, builtin origin, memory/query target,
1264    /// non-destructive, and engine-side shape verification: SUPERSEDE-only
1265    /// structural curation (never an ADD that introduces evidence-derived
1266    /// text) whose every replacement is **value-identical** to the grain it
1267    /// supersedes (the exact-equality check — a near-duplicate consolidation
1268    /// stays pending). A default (closed) policy never grants, so nothing
1269    /// auto-applies.
1270    fn can_auto_apply<S: OmsSubstrate>(&self, sub: &S, rec: &Recommendation) -> bool {
1271        if !rec.origin.auto_apply_eligible() || rec.destructive {
1272            return false;
1273        }
1274        // The analyzer must declare its curation auto-appliable. An analyzer
1275        // whose manifest is `Never` (e.g. fork surfacing — a lossy merge) is
1276        // never auto-applied even if the payload passes the shape check.
1277        let manifest_ok = self
1278            .analyzers
1279            .iter()
1280            .map(|a| a.manifest())
1281            .find(|m| m.id == rec.analyzer)
1282            .is_some_and(|m| m.auto_apply == crate::manifest::AutoApplyClass::StructuralCuration);
1283        if !manifest_ok {
1284            return false;
1285        }
1286        let Ok(target) = TargetRef::parse(&rec.target_ref) else {
1287            return false;
1288        };
1289        let family = crate::manifest::analyzer_family(&rec.analyzer);
1290        if !self.policy.grants_auto_apply(family, target.target_class(), rec.severity) {
1291            return false;
1292        }
1293        // Shape verification: only a CAL batch of pure SUPERSEDE statements
1294        // whose replacements change no value is structural curation. An ADD
1295        // (introducing content), a FORGET (destructive), or a supersession
1296        // that alters any field disqualifies.
1297        match &rec.proposal {
1298            Proposal::Cal { cal } => cal
1299                .lines()
1300                .map(str::trim)
1301                .filter(|l| !l.is_empty())
1302                .all(|l| supersede_is_value_identical(sub, l)),
1303            _ => false,
1304        }
1305    }
1306
1307    /// Apply a recommendation as `policy:auto` (the only `pending → applied`
1308    /// path). Records the applied inverse + a hash-chained audit grain.
1309    fn auto_apply<S: OmsSubstrate>(
1310        &self,
1311        sub: &mut S,
1312        p: &mut LoopPersisted,
1313        rec: &Recommendation,
1314        now_ms: i64,
1315    ) -> Result<()> {
1316        let mut created = Vec::new();
1317        if let Proposal::Cal { cal } = &rec.proposal {
1318            // Belt and braces over the policy: `grants_auto_apply` already
1319            // excludes the `query` class, so a definition rewrite cannot reach
1320            // this path. If one ever did, it would apply with no recorded
1321            // inverse and no human BECAUSE — refuse instead.
1322            if cal.lines().map(str::trim).any(is_definition_statement) {
1323                return Err(Error::InvalidProposal(
1324                    "a definition rewrite (DEFINE QUERY / DEFINE TEMPLATE) is never \
1325                     auto-applied: it changes what every future context contains, so it \
1326                     requires a human APPROVE + APPLY with BECAUSE"
1327                        .into(),
1328                ));
1329            }
1330            for r in sub.execute_cal(cal)? {
1331                if let Some(h) = r.get("hash").and_then(Value::as_str) {
1332                    created.push(h.to_string());
1333                }
1334            }
1335        }
1336        let applied = AppliedRecord {
1337            applied_at_ms: now_ms,
1338            target_ref: rec.target_ref.clone(),
1339            rollbackable: rec.rollbackable,
1340            created_hashes: created,
1341            inverse_cal: None,
1342            metric: rec.metric.clone(),
1343        };
1344        let prev = p.audit_heads.get(&rec.hash).cloned();
1345        let audit = AuditRecord {
1346            rec_hash: rec.hash.clone(),
1347            from: Some(RecStatus::Pending),
1348            to: RecStatus::Applied,
1349            actor: "policy:auto".into(),
1350            observer_type: ObserverType::Policy,
1351            because: "auto-applied per host policy".into(),
1352            previous_audit_hash: prev,
1353            gating: None,
1354            at_ms: now_ms,
1355        };
1356        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1357        p.audit_heads.insert(rec.hash.clone(), audit_hash);
1358        p.status_index.insert(rec.hash.clone(), RecStatus::Applied);
1359        p.applied.insert(rec.hash.clone(), applied);
1360        Ok(())
1361    }
1362
1363    /// Approve or reject a pending recommendation. Requires the `review` scope,
1364    /// a mandatory BECAUSE, and blocks self-approval against the creating actor.
1365    #[allow(clippy::too_many_arguments)]
1366    pub fn review<S: OmsSubstrate>(
1367        &self,
1368        sub: &mut S,
1369        rec_hash: &str,
1370        decision: Decision,
1371        actor: &str,
1372        observer: ObserverType,
1373        scopes: &ScopeSet,
1374        because: &str,
1375        now_ms: i64,
1376    ) -> Result<()> {
1377        if !scopes.has(Scope::Review) {
1378            return Err(Error::ScopeDenied("review".into()));
1379        }
1380        let because = validate_because(because)?;
1381        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1382        let status = *p
1383            .status_index
1384            .get(rec_hash)
1385            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1386        let to = match decision {
1387            Decision::Approve => RecStatus::Approved,
1388            Decision::Reject => RecStatus::Rejected,
1389        };
1390        if !status.can_transition_to(to, false) {
1391            return Err(Error::LifecycleViolation(format!(
1392                "{} -> {}",
1393                status.as_str(),
1394                to.as_str()
1395            )));
1396        }
1397        if to == RecStatus::Approved {
1398            if let Some(creator) = p.creators.get(rec_hash) {
1399                if creator == actor {
1400                    return Err(Error::SelfApproval(format!(
1401                        "{actor} created this recommendation"
1402                    )));
1403                }
1404            }
1405            if let Some(trigger) = p.co_creators.get(rec_hash) {
1406                if trigger == actor {
1407                    return Err(Error::SelfApproval(format!(
1408                        "{actor} triggered the run that authored this recommendation"
1409                    )));
1410                }
1411            }
1412        }
1413        let prev = p.audit_heads.get(rec_hash).cloned();
1414        let audit = AuditRecord {
1415            rec_hash: rec_hash.into(),
1416            from: Some(status),
1417            to,
1418            actor: actor.into(),
1419            observer_type: observer,
1420            because,
1421            previous_audit_hash: prev,
1422            gating: None,
1423            at_ms: now_ms,
1424        };
1425        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1426        p.audit_heads.insert(rec_hash.into(), audit_hash);
1427        p.status_index.insert(rec_hash.into(), to);
1428        if to == RecStatus::Rejected {
1429            if let Ok(rec) = load_rec(sub, rec_hash) {
1430                strike_cooldown(&mut p, rec.dedup_key, now_ms);
1431            }
1432        }
1433        sub.store_state(&p.to_value()?)?;
1434        Ok(())
1435    }
1436
1437    /// Check everything [`apply`](Self::apply) would refuse on, without writing
1438    /// anything.
1439    ///
1440    /// This exists for the fused approve-and-apply callers (the bindings'
1441    /// `apply_recommendation`). Recording the approval first and *then* hitting
1442    /// the destructive gate strands the recommendation in `approved`, which has
1443    /// no exit but `applied` or `expired` — `approved → rejected` is not a
1444    /// legal transition — so a refused apply left the reviewer unable to
1445    /// dismiss it. Ask first, then approve.
1446    ///
1447    /// Deliberately does not check the lifecycle transition: the caller is
1448    /// about to make it legal by approving.
1449    /// `has_gating` is whether the caller will supply a gating run at apply:
1450    /// a gated revision (code or adapter) without one is refused HERE, before
1451    /// a fused approve-and-apply records the approval — `approved` has no
1452    /// exit but `applied` or `expired`, so asking after would strand it.
1453    pub fn preflight_apply<S: OmsSubstrate>(
1454        &self,
1455        sub: &S,
1456        rec_hash: &str,
1457        scopes: &ScopeSet,
1458        allow_destructive: bool,
1459        has_gating: bool,
1460    ) -> Result<()> {
1461        if !scopes.has(Scope::Apply) {
1462            return Err(Error::ScopeDenied("apply".into()));
1463        }
1464        let rec = load_rec(sub, rec_hash)?;
1465        if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1466            return Err(Error::DestructiveGated(
1467                "destructive apply requires admin scope + allow_destructive".into(),
1468            ));
1469        }
1470        ensure_executable(rec.action_kind, &rec.proposal)?;
1471        if requires_gating(rec.action_kind) && !has_gating {
1472            return Err(Error::InvalidProposal(GATING_REQUIRED.into()));
1473        }
1474        Ok(())
1475    }
1476
1477    /// Apply an approved recommendation. Requires `apply`; destructive payloads
1478    /// additionally require `admin` + `allow_destructive`. Records the applied
1479    /// info (inverse plan) for rollback.
1480    #[allow(clippy::too_many_arguments)]
1481    pub fn apply<S: OmsSubstrate>(
1482        &self,
1483        sub: &mut S,
1484        rec_hash: &str,
1485        actor: &str,
1486        observer: ObserverType,
1487        scopes: &ScopeSet,
1488        because: &str,
1489        allow_destructive: bool,
1490        now_ms: i64,
1491    ) -> Result<AppliedRecord> {
1492        self.apply_inner(
1493            sub, rec_hash, actor, observer, scopes, because, allow_destructive, None, now_ms,
1494        )
1495    }
1496
1497    /// Load the gating evidence for `rec_hash` from the RECORDED
1498    /// `mg:eval_run` summary named by `run_id` — the one loader every
1499    /// surface (CLI, bindings, MCP, HTTP) shares, so the stats that admit a
1500    /// gated revision can never come from a caller: they are read back from
1501    /// the Fact `areev eval run` journaled, within the recommendation's own
1502    /// pinned evalset (a run id from a different evalset simply isn't found).
1503    pub fn gating_evidence<S: OmsSubstrate>(
1504        &self,
1505        sub: &S,
1506        rec_hash: &str,
1507        run_id: &str,
1508    ) -> Result<crate::recommendation::GatingEvidence> {
1509        let rec = self
1510            .recommendations(sub, None)?
1511            .into_iter()
1512            .find(|r| r.hash == rec_hash)
1513            .ok_or_else(|| {
1514                Error::InvalidProposal(format!("recommendation {rec_hash} not found"))
1515            })?;
1516        let pin = rec.evalset_hash.ok_or_else(|| {
1517            Error::InvalidProposal(
1518                "this recommendation pins no evalset — a gating run applies only \
1519                 to code and adapter revisions"
1520                    .into(),
1521            )
1522        })?;
1523        // Read through the shared evalset reader, so the run that GATES an
1524        // apply and the runs that later JUDGE it are parsed by exactly one
1525        // piece of code — two parsers drifting apart would let a rule be
1526        // admitted on one reading of a summary and measured on another.
1527        match crate::eval::eval_run_by_id(sub, &pin, run_id)? {
1528            Some(run) => Ok(crate::recommendation::GatingEvidence {
1529                evalset_hash: pin,
1530                run_id: run.run_id,
1531                passed: run.passed,
1532                failed: run.failed,
1533            }),
1534            None => Err(Error::InvalidProposal(format!(
1535                "no recorded gate run '{run_id}' for evalset {pin} — run \
1536                 `areev eval run --evalset {pin} ...` first"
1537            ))),
1538        }
1539    }
1540
1541    /// Apply WITH the §7.4 evalset-run edge — the only path that can apply a
1542    /// gated (code or adapter) revision. The evidence is validated against
1543    /// the recommendation's pin and recorded on the audit Observation.
1544    #[allow(clippy::too_many_arguments)]
1545    pub fn apply_gated<S: OmsSubstrate>(
1546        &self,
1547        sub: &mut S,
1548        rec_hash: &str,
1549        actor: &str,
1550        observer: ObserverType,
1551        scopes: &ScopeSet,
1552        because: &str,
1553        allow_destructive: bool,
1554        gating: &crate::recommendation::GatingEvidence,
1555        now_ms: i64,
1556    ) -> Result<AppliedRecord> {
1557        self.apply_inner(
1558            sub,
1559            rec_hash,
1560            actor,
1561            observer,
1562            scopes,
1563            because,
1564            allow_destructive,
1565            Some(gating),
1566            now_ms,
1567        )
1568    }
1569
1570    #[allow(clippy::too_many_arguments)]
1571    fn apply_inner<S: OmsSubstrate>(
1572        &self,
1573        sub: &mut S,
1574        rec_hash: &str,
1575        actor: &str,
1576        observer: ObserverType,
1577        scopes: &ScopeSet,
1578        because: &str,
1579        allow_destructive: bool,
1580        gating: Option<&crate::recommendation::GatingEvidence>,
1581        now_ms: i64,
1582    ) -> Result<AppliedRecord> {
1583        if !scopes.has(Scope::Apply) {
1584            return Err(Error::ScopeDenied("apply".into()));
1585        }
1586        let because = validate_because(because)?;
1587        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1588        let status = *p
1589            .status_index
1590            .get(rec_hash)
1591            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1592        if !status.can_transition_to(RecStatus::Applied, false) {
1593            return Err(Error::LifecycleViolation(format!(
1594                "{} -> applied (approve first)",
1595                status.as_str()
1596            )));
1597        }
1598        let rec = load_rec(sub, rec_hash)?;
1599        if rec.destructive && (!scopes.has(Scope::Admin) || !allow_destructive) {
1600            return Err(Error::DestructiveGated(
1601                "destructive apply requires admin scope + allow_destructive".into(),
1602            ));
1603        }
1604        // §7.4: a code or adapter revision applies ONLY through the
1605        // evalset-run edge — the pin must match, the pinned evalset must
1606        // still be LIVE (a superseded evalset invalidates in-flight
1607        // recommendations: they re-gate), and a failing gate admits nothing.
1608        if requires_gating(rec.action_kind) {
1609            let g = gating
1610                .ok_or_else(|| Error::InvalidProposal(GATING_REQUIRED.into()))?;
1611            let pin = rec.evalset_hash.as_deref().unwrap_or_default();
1612            if g.evalset_hash != pin {
1613                return Err(Error::InvalidProposal(format!(
1614                    "gating ran evalset {} but the recommendation is pinned \
1615                     to {pin} (Rule E1)",
1616                    g.evalset_hash
1617                )));
1618            }
1619            match sub.grain(pin)? {
1620                Some(evalset) if evalset.is_live() => {}
1621                Some(_) => {
1622                    return Err(Error::InvalidProposal(
1623                        "the pinned evalset was superseded after gating — \
1624                         the recommendation must re-gate (Rule E1)"
1625                            .into(),
1626                    ))
1627                }
1628                None => {
1629                    return Err(Error::InvalidProposal(format!(
1630                        "pinned evalset {pin} not found in the substrate"
1631                    )))
1632                }
1633            }
1634            if g.failed > 0 {
1635                return Err(Error::InvalidProposal(format!(
1636                    "the gating run failed {}/{} cases — a failing gate \
1637                     admits nothing",
1638                    g.failed,
1639                    g.passed + g.failed
1640                )));
1641            }
1642        }
1643
1644        // Execute the proposal.
1645        let mut created = Vec::new();
1646        // The inverse of a change that creates no grain — see
1647        // `AppliedRecord::inverse_cal`. Captured BEFORE execution, because
1648        // afterwards the previous definition is gone.
1649        let mut inverse_cal: Option<String> = None;
1650        match &rec.proposal {
1651            Proposal::Cal { cal } => {
1652                for line in cal.lines().map(str::trim).filter(|l| !l.is_empty()) {
1653                    if !is_definition_statement(line) {
1654                        continue;
1655                    }
1656                    match sub.definition_inverse(line)? {
1657                        Some(inv) => inverse_cal = Some(inv),
1658                        None => {
1659                            return Err(Error::InvalidProposal(format!(
1660                                "this substrate cannot record a rollback inverse for {line:?}; \
1661                                 a definition rewrite that ROLLBACK could not undo is refused \
1662                                 rather than applied"
1663                            )))
1664                        }
1665                    }
1666                }
1667                let rows = sub.execute_cal(cal)?;
1668                for r in rows {
1669                    if let Some(h) = r.get("hash").and_then(Value::as_str) {
1670                        created.push(h.to_string());
1671                    }
1672                }
1673            }
1674            // The engine has no executable Edit primitive. Marking this
1675            // Applied used to be a lie (and rollback had no inverse).
1676            Proposal::Edit { .. } => {
1677                return Err(Error::InvalidProposal(ADVISORY_EDIT.into()));
1678            }
1679            // A gated code or adapter revision EXECUTES by writing the
1680            // promotion grain: an immutable record that this target now
1681            // resolves to the proposed code (a §7.4 blob address) or adapter
1682            // (the tuning seam's registry tuple). Hosts read the promotion
1683            // to re-resolve — `(tool:X, mg:code_promotion)` or
1684            // `(model:X, mg:adapter_promotion)`; retracting it is the
1685            // rollback inverse, so the apply is rollbackable end-to-end.
1686            Proposal::Data { data } if requires_gating(rec.action_kind) => {
1687                let relation = if rec.action_kind == ActionKind::AdapterRevision {
1688                    "mg:adapter_promotion"
1689                } else {
1690                    "mg:code_promotion"
1691                };
1692                // A code revision authored by DISCOVER carries its SOURCE
1693                // inline, because the discovery pass reads the substrate and
1694                // cannot write to it. Move it into the CAS here so the
1695                // promotion names a content ADDRESS: §7.4's rule that code
1696                // enters the substrate only through the blob seam holds
1697                // however the revision was authored, and the promotion grain
1698                // stays a pointer rather than swelling to hold a program.
1699                let mut promoted = data.clone();
1700                if let Some(Value::String(src)) = promoted.remove("source") {
1701                    let address = sub.put_blob(src.as_bytes())?;
1702                    promoted.insert("code_address".into(), Value::from(address));
1703                }
1704                let mut spec = crate::substrate::GrainSpec::new(
1705                    crate::model::grain_type::FACT,
1706                    LOOP_NS,
1707                )
1708                .with_field("subject", rec.target_ref.clone())
1709                .with_field("relation", relation)
1710                .with_field(
1711                    "object",
1712                    serde_json::to_string(&promoted)
1713                        .map_err(|e| Error::Internal(format!("encode promotion: {e}")))?,
1714                )
1715                .with_field("rec_hash", rec_hash.to_string());
1716                if let Some(g) = gating {
1717                    spec = spec
1718                        .with_field("gating_evalset", g.evalset_hash.clone())
1719                        .with_field("gating_run_id", g.run_id.clone());
1720                }
1721                created.push(sub.put_grain(&spec)?);
1722            }
1723            Proposal::Data { data } => {
1724                // OutcomeReview is the one executable Data shape: its
1725                // `revert_of` points at an earlier applied recommendation.
1726                // Reuse the ordinary rollback path so the created hashes are
1727                // really retracted and the original lifecycle/audit advances.
1728                let revert_of = data
1729                    .get("revert_of")
1730                    .and_then(Value::as_str)
1731                    .ok_or_else(|| Error::InvalidProposal(ADVISORY_DATA.into()))?;
1732                self.rollback(
1733                    sub,
1734                    revert_of,
1735                    actor,
1736                    observer,
1737                    scopes,
1738                    &because,
1739                    now_ms,
1740                )?;
1741                // rollback stored a newer lifecycle state; merge this apply
1742                // into that state rather than overwriting the rollback.
1743                p = LoopPersisted::from_value(sub.load_state()?)?;
1744                // A measured revert is a verdict on the finding, not only on
1745                // this apply: the lesson was tried and it hurt. Rolled-back
1746                // findings normally re-propose ("the situation returned"),
1747                // which is right for an operator's rollback — but here the
1748                // situation never left, so the next pass would re-propose the
1749                // same lesson at once and the reviewer would be asked to
1750                // re-approve what the Verify gate just retracted. Put the
1751                // reverted finding on the same doubling cooldown a rejection
1752                // earns; the operator can still re-propose it by hand.
1753                if let Ok(reverted) = load_rec(sub, revert_of) {
1754                    strike_cooldown(&mut p, reverted.dedup_key, now_ms);
1755                }
1756            }
1757        }
1758
1759        let applied = AppliedRecord {
1760            applied_at_ms: now_ms,
1761            target_ref: rec.target_ref.clone(),
1762            rollbackable: rec.rollbackable,
1763            created_hashes: created,
1764            inverse_cal,
1765            metric: rec.metric.clone(),
1766        };
1767        let prev = p.audit_heads.get(rec_hash).cloned();
1768        let audit = AuditRecord {
1769            rec_hash: rec_hash.into(),
1770            from: Some(status),
1771            to: RecStatus::Applied,
1772            actor: actor.into(),
1773            observer_type: observer,
1774            because,
1775            previous_audit_hash: prev,
1776            gating: gating.cloned(),
1777            at_ms: now_ms,
1778        };
1779        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1780        p.audit_heads.insert(rec_hash.into(), audit_hash);
1781        p.status_index.insert(rec_hash.into(), RecStatus::Applied);
1782        p.applied.insert(rec_hash.into(), applied.clone());
1783        sub.store_state(&p.to_value()?)?;
1784        Ok(applied)
1785    }
1786
1787    /// Roll back an applied recommendation by retracting the grains it created.
1788    /// Fails for non-rollbackable applies (e.g. FORGET).
1789    #[allow(clippy::too_many_arguments)]
1790    pub fn rollback<S: OmsSubstrate>(
1791        &self,
1792        sub: &mut S,
1793        rec_hash: &str,
1794        actor: &str,
1795        observer: ObserverType,
1796        scopes: &ScopeSet,
1797        because: &str,
1798        now_ms: i64,
1799    ) -> Result<()> {
1800        if !scopes.has(Scope::Apply) {
1801            return Err(Error::ScopeDenied("apply".into()));
1802        }
1803        let because = validate_because(because)?;
1804        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1805        let status = *p
1806            .status_index
1807            .get(rec_hash)
1808            .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
1809        if !status.can_transition_to(RecStatus::RolledBack, false) {
1810            return Err(Error::LifecycleViolation(format!(
1811                "{} -> rolled_back",
1812                status.as_str()
1813            )));
1814        }
1815        let applied = p
1816            .applied
1817            .get(rec_hash)
1818            .cloned()
1819            .ok_or_else(|| Error::NotFound(format!("no applied record for {rec_hash}")))?;
1820        if !applied.rollbackable {
1821            return Err(Error::LifecycleViolation(
1822                "recommendation is non-rollbackable (FORGET has no inverse)".into(),
1823            ));
1824        }
1825        for h in &applied.created_hashes {
1826            sub.retract(h, &format!("rollback of {rec_hash}"))?;
1827        }
1828        // A definition rewrite creates no grain, so retracting `created_hashes`
1829        // undoes nothing. Restoring it means re-running the statement captured
1830        // at apply time — the previous definition, or a DROP when there was
1831        // none. Runs BEFORE the audit is written, so a failed restore leaves
1832        // the recommendation `applied` (still true) rather than recording a
1833        // rollback that did not happen.
1834        if let Some(inverse) = &applied.inverse_cal {
1835            sub.execute_cal(inverse)?;
1836        }
1837        let prev = p.audit_heads.get(rec_hash).cloned();
1838        let audit = AuditRecord {
1839            rec_hash: rec_hash.into(),
1840            from: Some(status),
1841            to: RecStatus::RolledBack,
1842            actor: actor.into(),
1843            observer_type: observer,
1844            because,
1845            previous_audit_hash: prev,
1846            gating: None,
1847            at_ms: now_ms,
1848        };
1849        let audit_hash = sub.put_grain(&audit.to_grain_spec(LOOP_NS))?;
1850        p.audit_heads.insert(rec_hash.into(), audit_hash);
1851        p.status_index
1852            .insert(rec_hash.into(), RecStatus::RolledBack);
1853        sub.store_state(&p.to_value()?)?;
1854        Ok(())
1855    }
1856
1857    /// List stored recommendations, optionally filtered by status. Status comes
1858    /// from the rebuildable index, not the immutable grain body. Ordered for
1859    /// review triage — highest severity first, then oldest first — and stable
1860    /// across runs for identical input.
1861    pub fn recommendations<S: OmsSubstrate>(
1862        &self,
1863        sub: &S,
1864        status_filter: Option<RecStatus>,
1865    ) -> Result<Vec<Recommendation>> {
1866        let p = LoopPersisted::from_value(sub.load_state()?)?;
1867        let grains = sub.grains_of_type(
1868            crate::model::grain_type::RECOMMENDATION,
1869            Some(LOOP_NS),
1870            ReadOpts {
1871                live_only: false,
1872                since_ms: None,
1873            },
1874        )?;
1875        let mut out = Vec::new();
1876        for g in grains {
1877            let mut rec = Recommendation::from_fields(&g.hash, &g.fields)?;
1878            rec.status = p
1879                .status_index
1880                .get(&g.hash)
1881                .copied()
1882                .unwrap_or(RecStatus::Pending);
1883            if let Some(f) = status_filter {
1884                if rec.status != f {
1885                    continue;
1886                }
1887            }
1888            out.push(rec);
1889        }
1890        // Review-queue order: worst first, then oldest first. Hash is only the
1891        // final tiebreak — sorting by it alone is deterministic per run but
1892        // meaningless across runs, because a grain's hash covers its timestamp,
1893        // so an identical queue comes back in a different order every time.
1894        // `dedup_key` is the last tiebreak that actually decides anything: it is
1895        // content-derived and stable across runs, whereas findings proposed in
1896        // the same sweep routinely share a `created_at_ms`. Hash trails it only
1897        // to make the ordering total.
1898        out.sort_by(|a, b| {
1899            b.severity
1900                .cmp(&a.severity)
1901                .then(a.created_at_ms.cmp(&b.created_at_ms))
1902                .then(a.dedup_key.cmp(&b.dedup_key))
1903                .then(a.hash.cmp(&b.hash))
1904        });
1905        Ok(out)
1906    }
1907
1908    /// Per-analyzer effective settings for the Setup view: the manifest facts
1909    /// merged with the file-config (override or manifest default). Read-only.
1910    pub fn analyzer_settings<S: OmsSubstrate>(
1911        &self,
1912        sub: &S,
1913    ) -> Result<Vec<crate::config::AnalyzerSetting>> {
1914        let p = LoopPersisted::from_value(sub.load_state()?)?;
1915        Ok(self
1916            .analyzers
1917            .iter()
1918            .map(|a| {
1919                let m = a.manifest();
1920                let cfg = p.config.get(&m.id);
1921                crate::config::AnalyzerSetting {
1922                    id: m.id.clone(),
1923                    title: m.title.clone(),
1924                    description: m.description.clone(),
1925                    tier: format!("{:?}", m.tier),
1926                    trust_class: format!("{:?}", m.trust_class).to_lowercase(),
1927                    default_on: m.default_on,
1928                    enabled: cfg.and_then(|c| c.enabled).unwrap_or(m.default_on),
1929                    severity_floor: cfg
1930                        .and_then(|c| c.severity_floor)
1931                        .map(|s| s.as_str().to_string()),
1932                }
1933            })
1934            .collect())
1935    }
1936
1937    /// Update one analyzer's file-config (enable/disable, severity floor, param
1938    /// overrides, namespace scoping). Requires `Admin`. Params are validated
1939    /// against the analyzer's manifest first (unknown keys rejected, fail-closed),
1940    /// and the analyzer must exist. Returns the merged config as stored. This is
1941    /// the only write into `persisted.config` — the config layer, never a grain.
1942    pub fn set_analyzer_config<S: OmsSubstrate>(
1943        &self,
1944        sub: &mut S,
1945        analyzer_id: &str,
1946        update: crate::config::AnalyzerConfigUpdate,
1947        scopes: &ScopeSet,
1948    ) -> Result<crate::config::AnalyzerConfig> {
1949        if !scopes.has(Scope::Admin) {
1950            return Err(Error::ScopeDenied("admin".into()));
1951        }
1952        let manifest = self
1953            .analyzers
1954            .iter()
1955            .map(|a| a.manifest())
1956            .find(|m| m.id == analyzer_id)
1957            .ok_or_else(|| Error::NotFound(format!("unknown analyzer {analyzer_id:?}")))?;
1958        // Validate params against the manifest BEFORE touching state.
1959        if let Some(params) = &update.params {
1960            manifest.resolve_params(params)?;
1961        }
1962        let mut p = LoopPersisted::from_value(sub.load_state()?)?;
1963        let cfg = p.config.entry(analyzer_id.to_string()).or_default();
1964        if let Some(enabled) = update.enabled {
1965            cfg.enabled = Some(enabled);
1966        }
1967        if update.clear_floor {
1968            cfg.severity_floor = None;
1969        } else if let Some(floor) = update.severity_floor {
1970            cfg.severity_floor = Some(floor);
1971        }
1972        if let Some(params) = update.params {
1973            cfg.params = params;
1974        }
1975        if let Some(ns) = update.namespaces {
1976            cfg.namespaces = ns;
1977        }
1978        let stored = cfg.clone();
1979        sub.store_state(&p.to_value()?)?;
1980        Ok(stored)
1981    }
1982
1983    /// The measured outcome time series (the Verify gate's history) across all
1984    /// recommendations, ordered by when each checkpoint was measured.
1985    pub fn outcomes<S: OmsSubstrate>(&self, sub: &S) -> Result<Vec<crate::recommendation::OutcomeResult>> {
1986        let p = LoopPersisted::from_value(sub.load_state()?)?;
1987        let mut out: Vec<_> = p.outcomes.into_values().flatten().collect();
1988        // `metric` and `rec_hash` break the tie: checkpoints measured in the
1989        // same sweep share a `measured_at_ms`, and without a tiebreak the order
1990        // falls through to the map's rec_hash ordering, which shifts every run.
1991        out.sort_by(|a, b| {
1992            a.measured_at_ms
1993                .cmp(&b.measured_at_ms)
1994                .then(a.horizon_ms.cmp(&b.horizon_ms))
1995                .then(a.metric.cmp(&b.metric))
1996                .then(a.rec_hash.cmp(&b.rec_hash))
1997        });
1998        Ok(out)
1999    }
2000
2001    /// A health snapshot — when the loop last ran, how much is un-analyzed
2002    /// since, and the queue counts. Lets a host surface "the loop may be stale"
2003    /// so a forgotten SessionEnd hook / cron doesn't silently kill it.
2004    pub fn health<S: OmsSubstrate>(&self, sub: &S, now_ms: i64) -> Result<Health> {
2005        let p = LoopPersisted::from_value(sub.load_state()?)?;
2006        let new = count_new(sub, p.state.watermark_ms)?;
2007        let (grains_since_run, error_events_since_run) = (new.grains, new.error_events);
2008        let recs = self.recommendations(sub, None)?;
2009        let mut pending = 0;
2010        let mut applied = 0;
2011        for r in &recs {
2012            match r.status {
2013                RecStatus::Pending => pending += 1,
2014                RecStatus::Applied => applied += 1,
2015                _ => {}
2016            }
2017        }
2018        // Stale if it has never run, or it's been a while / a lot has piled up.
2019        let stale = match p.state.last_run_ms {
2020            None => true,
2021            Some(last) => now_ms - last >= 7 * 86_400_000 || grains_since_run >= 100,
2022        };
2023        Ok(Health {
2024            last_run_ms: p.state.last_run_ms,
2025            grains_since_run,
2026            error_events_since_run,
2027            pending,
2028            applied,
2029            total: recs.len() as u64,
2030            stale,
2031        })
2032    }
2033
2034    /// Approval-rate metric for `origin = llm` recommendations (reflection
2035    /// design §6b) — the live field-quality signal that accrues off the audit
2036    /// chain: what fraction of the model's *surfaced* proposals a reviewer
2037    /// accepts. Complements the offline Effective-Reliability eval.
2038    pub fn llm_metrics<S: OmsSubstrate>(&self, sub: &S) -> Result<LlmMetrics> {
2039        let recs = self.recommendations(sub, None)?;
2040        let mut m = LlmMetrics::default();
2041        for r in &recs {
2042            if !matches!(r.origin, Origin::Llm { .. }) {
2043                continue;
2044            }
2045            m.proposed += 1;
2046            match r.status {
2047                RecStatus::Pending => m.pending += 1,
2048                RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack => m.approved += 1,
2049                RecStatus::Rejected => m.rejected += 1,
2050                RecStatus::Expired => {}
2051            }
2052        }
2053        let decided = m.approved + m.rejected;
2054        m.approval_rate = (decided > 0).then(|| m.approved as f64 / decided as f64);
2055        Ok(m)
2056    }
2057}
2058
2059/// A health snapshot for the backend's self-improvement loop.
2060#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
2061pub struct Health {
2062    #[serde(skip_serializing_if = "Option::is_none")]
2063    pub last_run_ms: Option<i64>,
2064    pub grains_since_run: u64,
2065    pub error_events_since_run: u64,
2066    pub pending: u64,
2067    pub applied: u64,
2068    pub total: u64,
2069    /// True when the loop looks stalled (never run, or ≥7d / ≥100 new grains
2070    /// since the last run) — a nudge that a trigger may be unwired.
2071    pub stale: bool,
2072}
2073
2074/// Approval-rate metric for `origin = llm` recommendations (reflection §6b).
2075#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
2076pub struct LlmMetrics {
2077    /// Total llm-origin recommendations ever stored (those that survived the
2078    /// verifier and reached the queue).
2079    pub proposed: u64,
2080    pub pending: u64,
2081    /// Approved + Applied + RolledBack (a reviewer said yes at least once).
2082    pub approved: u64,
2083    pub rejected: u64,
2084    /// approved / (approved + rejected); `None` until at least one is decided.
2085    #[serde(skip_serializing_if = "Option::is_none")]
2086    pub approval_rate: Option<f64>,
2087}
2088
2089/// Re-measure applied recommendations at each **checkpoint** past due, via the
2090/// engine's typed reads — no CAL-scalar round-trip. A recommendation
2091/// accumulates one `OutcomeResult` per horizon (measured once each), forming a
2092/// time series, so a late regression (held at 1d, regressed at 30d) is caught.
2093/// Only *regressed* checkpoints feed the outcome analyzer (→ a revert).
2094/// Unknown metric kinds are skipped, never faked.
2095fn measure_outcomes<S: OmsSubstrate>(
2096    sub: &S,
2097    p: &mut LoopPersisted,
2098    policy: &crate::policy::Policy,
2099    now_ms: i64,
2100) -> Result<Vec<OutcomeInput>> {
2101    // Collect all due (recommendation, checkpoint) pairs first. A checkpoint
2102    // is due in its own unit: elapsed time, evalset runs journaled since the
2103    // apply, or grains written since it (`Checkpoint`).
2104    let mut due: Vec<(String, crate::config::AppliedRecord, Checkpoint)> = Vec::new();
2105    for (h, a) in &p.applied {
2106        if p.status_index.get(h) != Some(&RecStatus::Applied) {
2107            continue;
2108        }
2109        let Some(metric) = &a.metric else { continue };
2110        let done = p.measured.get(h).cloned().unwrap_or_default();
2111        for cp in metric.schedule() {
2112            if !done.contains(&cp) && checkpoint_due(sub, metric, a.applied_at_ms, cp, now_ms)? {
2113                due.push((h.clone(), a.clone(), cp));
2114            }
2115        }
2116    }
2117
2118    let mut out = Vec::new();
2119    for (rec_hash, applied, checkpoint) in due {
2120        let metric = applied.metric.as_ref().unwrap();
2121        let Some(current) = measure_metric(sub, metric, applied.applied_at_ms)? else {
2122            continue; // metric kind not yet re-measurable
2123        };
2124        let bound = cost_bound_for(policy, metric);
2125        let base = baseline_at_apply(
2126            sub,
2127            metric,
2128            applied.applied_at_ms,
2129            baseline_kind_for(policy, metric),
2130            bound.map(|b| b.field.as_str()),
2131        )?;
2132        let baseline = base.value;
2133        // The floor is resolved ONCE here and travels with the input, so the
2134        // recorded verdict and the revert draft cannot disagree on it.
2135        let tolerance = tolerance_for(policy, metric, &base);
2136        let regressed = crate::recommendation::is_regression(
2137            baseline,
2138            current,
2139            metric.higher_is_better,
2140            tolerance,
2141        );
2142        // The run `current` came from, and the cost bound's reading on it.
2143        // The quality verdict never depends on the cost: `regressed`
2144        // dominates, a breached bound under a held score is `held_costlier`,
2145        // and a cost that is not measurable leaves the verdict alone.
2146        let (current_run_id, cost) = current_run_and_cost(sub, metric, applied.applied_at_ms, bound, &base)?;
2147        let costlier = cost.as_ref().is_some_and(|c| c.breached());
2148        let verdict = if regressed {
2149            "regressed"
2150        } else if costlier {
2151            "held_costlier"
2152        } else {
2153            "held"
2154        };
2155        p.outcomes.entry(rec_hash.clone()).or_default().push(
2156            crate::recommendation::OutcomeResult {
2157                rec_hash: rec_hash.clone(),
2158                metric: metric.metric.clone(),
2159                baseline,
2160                current,
2161                verdict: verdict.into(),
2162                baseline_kind: base.kind.into(),
2163                baseline_run_id: base.run_id.clone(),
2164                best_before: base.best_before,
2165                tolerance,
2166                current_run_id: current_run_id.clone(),
2167                cost: cost.clone(),
2168                // A time checkpoint speaks through `horizon_ms`, as it always
2169                // did; the other units carry themselves.
2170                horizon_ms: checkpoint.as_ms().unwrap_or(0),
2171                checkpoint: checkpoint.as_ms().is_none().then_some(checkpoint),
2172                measured_at_ms: now_ms,
2173            },
2174        );
2175        p.measured.entry(rec_hash.clone()).or_default().push(checkpoint);
2176        if regressed || costlier {
2177            out.push(OutcomeInput {
2178                rec_hash,
2179                target_ref: applied.target_ref.clone(),
2180                metric: metric.metric.clone(),
2181                baseline,
2182                current,
2183                unit: metric.unit.clone(),
2184                higher_is_better: metric.higher_is_better,
2185                baseline_kind: base.kind.into(),
2186                baseline_run_id: base.run_id,
2187                best_before: base.best_before,
2188                tolerance,
2189                current_run_id,
2190                cost,
2191            });
2192        }
2193    }
2194    Ok(out)
2195}
2196
2197/// The policy's cost bound, for the evalset the policy names.
2198fn cost_bound_for<'p>(
2199    policy: &'p crate::policy::Policy,
2200    metric: &crate::recommendation::MetricSnapshot,
2201) -> Option<&'p crate::policy::CostBound> {
2202    match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2203        (Some(e), Some((hash, _))) if e.hash == hash => e.cost.as_ref(),
2204        _ => None,
2205    }
2206}
2207
2208/// For an evalset metric: the run the current value was read from (the
2209/// same lookup `measure_metric` made — one reader, one run), and the cost
2210/// bound's reading of it against the baseline run. `not_measurable` when
2211/// either side lacks the field or carries it malformed; the figures that
2212/// did measure still ride on the record.
2213fn current_run_and_cost<S: SubstrateRead>(
2214    sub: &S,
2215    metric: &crate::recommendation::MetricSnapshot,
2216    applied_at_ms: i64,
2217    bound: Option<&crate::policy::CostBound>,
2218    base: &BaselineRead,
2219) -> Result<(Option<String>, Option<crate::recommendation::CostRead>)> {
2220    let Some((evalset, _)) = crate::eval::parse_evalset_metric(&metric.metric) else {
2221        return Ok((None, None));
2222    };
2223    let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(applied_at_ms))? else {
2224        return Ok((None, None));
2225    };
2226    let cost = bound.map(|b| {
2227        let current = crate::eval::run_value(&run, &b.field);
2228        let status = match (base.cost, current) {
2229            (Some(bl), Some(cur)) if cur > bl * b.max_increase_ratio + 1e-9 => "breached",
2230            (Some(_), Some(_)) => "within",
2231            _ => "not_measurable",
2232        };
2233        crate::recommendation::CostRead {
2234            field: b.field.clone(),
2235            max_increase_ratio: b.max_increase_ratio,
2236            baseline: base.cost,
2237            current,
2238            status: status.into(),
2239        }
2240    });
2241    Ok((Some(run.run_id), cost))
2242}
2243
2244/// The policy's minimum effect size in the metric's unit, for the evalset
2245/// the policy names; zero for everything else. The `points` form scales a
2246/// count field by the baseline run's total — or, with no run before the
2247/// apply, by the total the proposal's snapshot recorded.
2248fn tolerance_for(
2249    policy: &crate::policy::Policy,
2250    metric: &crate::recommendation::MetricSnapshot,
2251    base: &BaselineRead,
2252) -> f64 {
2253    match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2254        (Some(e), Some((hash, field))) if e.hash == hash => e
2255            .min_effect
2256            .map(|m| m.resolve(field, base.total.unwrap_or(metric.n)))
2257            .unwrap_or(0.0),
2258        _ => 0.0,
2259    }
2260}
2261
2262/// Has this checkpoint come due for a recommendation applied at `applied_at_ms`?
2263///
2264/// Time is a subtraction. Runs are counted from the evalset the metric names
2265/// — journaled strictly after the apply, the same set the measurement itself
2266/// reads, so "due" and "measurable" can never disagree; on a metric that is
2267/// not evalset-backed a run checkpoint never fires, honestly, rather than
2268/// guessing what a run would be. Grains are the loop's own activity count
2269/// since the apply.
2270fn checkpoint_due<S: SubstrateRead>(
2271    sub: &S,
2272    metric: &crate::recommendation::MetricSnapshot,
2273    applied_at_ms: i64,
2274    cp: Checkpoint,
2275    now_ms: i64,
2276) -> Result<bool> {
2277    Ok(match cp {
2278        Checkpoint::AfterMs(ms) => now_ms - applied_at_ms >= ms,
2279        Checkpoint::AfterRuns(n) => match crate::eval::parse_evalset_metric(&metric.metric) {
2280            Some((evalset, _)) => {
2281                crate::eval::eval_runs(sub, evalset, Some(applied_at_ms + 1))?.len() >= n as usize
2282            }
2283            None => false,
2284        },
2285        Checkpoint::AfterGrains(n) => count_new(sub, Some(applied_at_ms))?.grains >= n as u64,
2286    })
2287}
2288
2289/// The number a verdict compares against, and where it came from.
2290pub(crate) struct BaselineRead {
2291    pub value: f64,
2292    /// `newest_before_apply`, `high_water` or `snapshot` — see
2293    /// `OutcomeResult::baseline_kind`.
2294    pub kind: &'static str,
2295    pub run_id: Option<String>,
2296    pub best_before: Option<f64>,
2297    /// Cases the baseline run graded, when it was a run.
2298    pub total: Option<u64>,
2299    /// The policy's cost field on the baseline run, when both exist and it
2300    /// is measurable there.
2301    pub cost: Option<f64>,
2302}
2303
2304/// The host's baseline choice applies to the evalset it names; any other
2305/// evalset-backed metric keeps the marginal comparison.
2306fn baseline_kind_for(
2307    policy: &crate::policy::Policy,
2308    metric: &crate::recommendation::MetricSnapshot,
2309) -> crate::policy::BaselineKind {
2310    match (policy.outcome_evalset.as_ref(), crate::eval::parse_evalset_metric(&metric.metric)) {
2311        (Some(e), Some((hash, _))) if e.hash == hash => e.baseline,
2312        _ => crate::policy::BaselineKind::default(),
2313    }
2314}
2315
2316/// The number a verdict compares against: for an evalset metric, a run
2317/// journaled BEFORE the apply when there is one — the newest by default, the
2318/// best under `high_water` — else the snapshot the proposal froze. "Did
2319/// applying it help" is a question about the state of the world at the
2320/// apply, not at the proposal — and a deployment that measures once, at day
2321/// one, and then approves its twentieth rule would otherwise read a rule
2322/// that cost twenty points as `held` against a baseline the first nineteen
2323/// had already left far behind. Found on a real corpus
2324/// (`crates/areev-bench/CURVE.md`: a rule that contradicted an earlier one
2325/// took the agent from 86% to 66% and measured as held against 26%). With
2326/// nothing journaled between proposal and apply the two are the same run,
2327/// so no verdict recorded before this changes.
2328///
2329/// `best_before` is read regardless of the choice: the marginal verdict is
2330/// blind to a fall from the peak by construction (`ADBUY.md`, seed 3: 35 →
2331/// 238 → 128 → 133 reads `held` against 35), so the receipt carries the
2332/// peak beside the baseline even when it is not the baseline.
2333pub(crate) fn baseline_at_apply<S: SubstrateRead>(
2334    sub: &S,
2335    metric: &crate::recommendation::MetricSnapshot,
2336    applied_at_ms: i64,
2337    kind: crate::policy::BaselineKind,
2338    cost_field: Option<&str>,
2339) -> Result<BaselineRead> {
2340    use crate::policy::BaselineKind;
2341    if let Some((evalset, field)) = crate::eval::parse_evalset_metric(&metric.metric) {
2342        // Oldest first, the field read through the one reader every consumer
2343        // uses; a run whose summary lacks the field is not a candidate.
2344        let before: Vec<(crate::eval::EvalRun, f64)> = crate::eval::eval_runs(sub, evalset, None)?
2345            .into_iter()
2346            .filter(|r| r.recorded_ms < applied_at_ms)
2347            .filter_map(|r| crate::eval::run_value(&r, field).map(|v| (r, v)))
2348            .collect();
2349        if let Some(newest) = before.last() {
2350            // The first run to attain the best value is the high-water mark:
2351            // a later tie did not raise it.
2352            let best = before
2353                .iter()
2354                .fold(None::<&(crate::eval::EvalRun, f64)>, |acc, r| match acc {
2355                    None => Some(r),
2356                    Some(b) => {
2357                        let better = if metric.higher_is_better { r.1 > b.1 } else { r.1 < b.1 };
2358                        Some(if better { r } else { b })
2359                    }
2360                })
2361                .expect("non-empty");
2362            let pick = match kind {
2363                BaselineKind::NewestBeforeApply => newest,
2364                BaselineKind::HighWater => best,
2365            };
2366            return Ok(BaselineRead {
2367                value: pick.1,
2368                kind: kind.as_str(),
2369                run_id: Some(pick.0.run_id.clone()),
2370                best_before: Some(best.1),
2371                total: Some(pick.0.total()),
2372                cost: cost_field.and_then(|f| crate::eval::run_value(&pick.0, f)),
2373            });
2374        }
2375    }
2376    Ok(BaselineRead { value: metric.baseline, kind: "snapshot", run_id: None, best_before: None, total: None, cost: None })
2377}
2378
2379/// Typed re-measurement for the fixed set of metric kinds the engine knows.
2380pub(crate) fn measure_metric<S: SubstrateRead>(
2381    sub: &S,
2382    metric: &crate::recommendation::MetricSnapshot,
2383    since_ms: i64,
2384) -> Result<Option<f64>> {
2385    match metric.metric.as_str() {
2386        // How many times did this tool fail again *with the same signature*
2387        // after the lesson was applied? Scoped to the signature (metric.relation)
2388        // so an unrelated later failure of the same tool is not read as a
2389        // regression of this specific lesson.
2390        "tool_error_recurrence" => {
2391            let Some(tool) = &metric.subject else { return Ok(None) };
2392            let tools = sub.grains_of_type(
2393                crate::model::grain_type::TOOL,
2394                None,
2395                ReadOpts { live_only: true, since_ms: Some(since_ms) },
2396            )?;
2397            let n = tools
2398                .iter()
2399                .filter(|t| t.tool_name() == Some(tool.as_str()) && t.is_error())
2400                .filter(|t| {
2401                    // No stored signature (legacy metric) → fall back to the
2402                    // whole-tool count so old recommendations still measure.
2403                    metric.relation.as_deref().is_none_or(|sig| {
2404                        crate::analyzers::tool_failure::normalize_signature(
2405                            t.tool_content().unwrap_or(""),
2406                        ) == sig
2407                    })
2408                })
2409                .count();
2410            Ok(Some(n as f64))
2411        }
2412        // After a resolve-to-latest, does the subject again hold more than one
2413        // live value under the functional relation? Live-state read (no since
2414        // filter): the excess beyond one distinct object is the regression.
2415        "contradiction_recurrence" => {
2416            let (Some(subject), Some(relation)) = (&metric.subject, &metric.relation) else {
2417                return Ok(None);
2418            };
2419            let facts = scoped_live_facts(sub, metric.namespace.as_deref(), subject)?;
2420            let distinct: BTreeSet<String> = facts
2421                .iter()
2422                .filter(|f| {
2423                    f.fact_relation()
2424                        .is_some_and(|r| normalize_ident(r) == *relation)
2425                })
2426                .filter_map(|f| f.fact_object().map(normalize_ident))
2427                .collect();
2428            Ok(Some(distinct.len().saturating_sub(1) as f64))
2429        }
2430        // `evalset:<hash>:<field>` — external correctness, measured the only
2431        // way that keeps the honesty rule ("Areev Loop improves the agent's
2432        // memory, not its outputs") intact: an evalset run is an INTERNAL,
2433        // BOUNDED, ATTRIBUTABLE measurement. The engine never runs the evalset;
2434        // it only reads summaries a host journaled with `areev eval run`.
2435        //
2436        // `since_ms` is the apply time, and it is load-bearing rather than an
2437        // optimization: a run journaled BEFORE the apply cannot be evidence of
2438        // what applying did. Comparing the baseline run against itself would
2439        // report "held" forever — a fabricated receipt, which is worse than no
2440        // receipt at all. No run since the apply → `None` → not yet
2441        // measurable, and the checkpoint stays due.
2442        m if m.starts_with("evalset:") => {
2443            let Some((evalset, field)) = crate::eval::parse_evalset_metric(m) else {
2444                return Ok(None);
2445            };
2446            let Some(run) = crate::eval::newest_eval_run(sub, evalset, Some(since_ms))? else {
2447                return Ok(None);
2448            };
2449            Ok(crate::eval::run_value(&run, field))
2450        }
2451        _ => Ok(None),
2452    }
2453}
2454
2455/// Live facts for one normalized (namespace?, subject) — the shared scope of
2456/// the fact-shaped recurrence metrics. `namespace: None` spans all namespaces.
2457fn scoped_live_facts<S: SubstrateRead>(
2458    sub: &S,
2459    namespace: Option<&str>,
2460    subject: &str,
2461) -> Result<Vec<GrainRecord>> {
2462    let facts = sub.grains_of_type(
2463        crate::model::grain_type::FACT,
2464        None,
2465        ReadOpts { live_only: true, since_ms: None },
2466    )?;
2467    Ok(facts
2468        .into_iter()
2469        .filter(|f| namespace.is_none_or(|ns| normalize_ident(&f.namespace) == ns))
2470        .filter(|f| {
2471            f.fact_subject()
2472                .is_some_and(|s| normalize_ident(s) == subject)
2473        })
2474        .collect())
2475}
2476
2477// --- free helpers ---
2478
2479/// The §6.3 exact-equality check: a SUPERSEDE is *value-identical* when every
2480/// replacement field equals the superseded grain's value — strings after
2481/// case-fold/trim (upstream NFC is an OMS invariant), `namespace` against the
2482/// grain's own namespace, everything else exactly. This is what makes an
2483/// auto-applied consolidation provably information-preserving; a
2484/// near-duplicate (an observation body off by one token) fails it and stays
2485/// pending for human review. Fails closed: an unrecognized line shape, an
2486/// empty replacement, a missing grain, or a field the original never had all
2487/// disqualify.
2488fn supersede_is_value_identical<S: SubstrateRead>(sub: &S, line: &str) -> bool {
2489    let Some((target, _gtype, fields)) = crate::cal::parse_own_supersede(line) else {
2490        return false;
2491    };
2492    if fields.is_empty() {
2493        return false;
2494    }
2495    let Ok(Some(grain)) = sub.grain(&target) else {
2496        return false;
2497    };
2498    // An expiry lives OUTSIDE `fields`, so the replacement (built from a fixed
2499    // field set) can never carry it — consolidating away a grain that has a
2500    // valid_to would silently drop the expiry, invisibly to the field-by-field
2501    // check below. Fail closed. (A dup that additionally carries an extra
2502    // *content* field the replacement omits is a real but subtler info-loss;
2503    // catching it soundly needs a full canonical-vs-extra comparison rather
2504    // than this replacement-scoped check, since a real fact's fields also carry
2505    // OMS metadata like `confidence` that consolidation legitimately keeps —
2506    // left as a follow-up so this narrow fix can't block valid consolidations.)
2507    if grain.valid_to_ms.is_some() {
2508        return false;
2509    }
2510    // Forward check: every replacement field equals the grain's value.
2511    fields.iter().all(|(k, v)| {
2512        if k == "namespace" {
2513            return v
2514                .as_str()
2515                .is_some_and(|s| normalize_ident(s) == normalize_ident(&grain.namespace));
2516        }
2517        match (v, grain.fields.get(k)) {
2518            (Value::String(a), Some(Value::String(b))) => {
2519                normalize_ident(a) == normalize_ident(b)
2520            }
2521            (a, Some(b)) => a == b,
2522            (_, None) => false,
2523        }
2524    })
2525}
2526
2527/// The DISCOVER evidence bundle's total size, and the reserved share each
2528/// source gets inside it.
2529///
2530/// Three sources feed the bundle and they answer different questions, so each
2531/// is budgeted rather than served first-come: the grains the deterministic
2532/// findings CITED (what clustering already caught), recent tool ERRORS (what
2533/// clustering could have caught but did not), and recent facts/observations —
2534/// the LLM's own lens, and the ONLY source that can carry a problem with no
2535/// error shape at all. `LENS_RESERVE` is what the last of those is guaranteed.
2536const EVIDENCE_CAP: usize = 64;
2537const CITED_SEED_CAP: usize = 24;
2538const TOOL_SEED_CAP: usize = 16;
2539/// Human-authored Observations get a small guaranteed share, taken before the
2540/// general top-up. Learning does not only come from what went wrong: a person
2541/// saying "from now on, do X" is a complete rule stated once, and no amount of
2542/// counting recovers it from a corpus that never surfaced it.
2543const NOTE_SEED_CAP: usize = 8;
2544/// Harness Observations the LLM lens is allowed to see, and how many.
2545///
2546/// An all-namespace scan deliberately hides every `agent:` namespace
2547/// (`areev-loop-adapter`): those hold the file's own grants and Tier-2 audit
2548/// records, and an analyzer that swept them as ordinary memory once proposed
2549/// tombstoning the grants — which locks every non-owner out of the file. That
2550/// exclusion stays exactly as it is.
2551///
2552/// But not everything the harness records is governance. A fold summary is the
2553/// agent's OWN account of what a long run had worked out, written when the
2554/// transcript outgrew the model's window — experience, in the harness
2555/// namespace only because it is evidence about a run rather than memory the
2556/// agent asserts. Invisible to the lens, it may as well not have been written.
2557///
2558/// So: an explicit, named list read from an EXPLICIT namespace (which is what
2559/// distinguishes it from a sweep), with a reserve of its own. Adding a kind
2560/// here is a one-line decision someone reviews — never a blanket un-hiding of
2561/// `agent:*`.
2562const HARNESS_EVIDENCE_KINDS: &[&str] = &["fold_summary"];
2563const HARNESS_SEED_CAP: usize = 6;
2564const LENS_RESERVE: usize = 24;
2565
2566/// The confidence floor (§5.4): a verified draft below this is dropped. The
2567/// verifier's calibrated confidence is the gate, not the proposer's self-report.
2568const MIN_LLM_CONFIDENCE: f64 = 0.75;
2569
2570/// The fixed DISCOVER instruction (§5.1), in two objectives that differ in
2571/// exactly one paragraph — the scoring rule — so the vocabulary, the cite
2572/// rule and the JSON contract cannot drift between them. Kept in its own
2573/// request field so it never interleaves with (attacker-influenced) evidence
2574/// text. The review-queue rule makes "nothing to report" a first-class,
2575/// zero-penalty answer — the structural antidote to over-generation; the
2576/// learner rule makes abstaining over evidence that plainly holds a lesson
2577/// cost the same as a wrong one. Which applies is host policy
2578/// (`Policy::discover_objective`), never the model's or the file's choice.
2579macro_rules! discover_instructions {
2580    ($scoring:literal) => {
2581        concat!(
2582            "You review an agent's memory for quality. \
2583Given deterministic findings and the evidence they cite, propose ADDITIONAL \
2584findings the deterministic checks would miss (e.g. a semantic contradiction, a \
2585stale assumption, a duplicated meaning, a recurring preventable mistake, a \
2586recurring cost or hand-off the agent's own setup could remove). \
2587The deterministic findings already cover what the ERROR TEXT says; restating \
2588one of them earns nothing. The evidence may also contain OUTCOME records — a \
2589run's observable shape together with whether it was accepted or rejected. A \
2590problem that raised no error at all is exactly the kind the deterministic \
2591checks cannot see, so compare the rejected outcomes against the accepted \
2592ones: a feature they share and the accepted ones lack is a candidate rule. \
2593Require at least two rejected outcomes before proposing one — a single \
2594rejection is an anecdote, not a pattern. ",
2595            $scoring,
2596            " The 'approved' and 'rejected' lists, when \
2597present, show findings this reviewer recently accepted or rejected — prefer the \
2598kind they accept and avoid the kind they reject. Every proposal MUST cite one \
2599or more evidence items from the bundle by their 'id' (or 'hash'), name a \
2600'target', and include your confidence 0.0-1.0. Return JSON: \
2601{\"recommendations\":[{\"summary\":\"...\",\
2602\"target\":\"...\",\"guidance\":\"...\",\"evidence\":[\"<id>\"],\
2603\"confidence\":0.0,\"proposal\":{...}}]}. \
2604OMIT 'proposal' for an advisory finding — one worth a human's attention that \
2605you are not asking to change anything. Include it ONLY when the evidence \
2606supports a specific change, choosing exactly one kind: \
2607(1) {\"kind\":\"lesson\",\"lesson\":\"...\"} with target \
2608\"entity:<ns>/<subject>\" — ONE short imperative rule (max 240 chars) naming \
2609an action the agent itself takes on the next occasion. Either ADD an action \
2610it is failing to take ('Record the vendor name and the amount on every \
2611invoice, not just the date') or ORDER one that goes wrong ('Refund a \
2612subscription before cancelling it; refunds on cancelled subscriptions are \
2613refused'). Name the action, not a check on it: 'validate', 'verify' and \
2614'ensure ... is correct' describe a review step the agent has no way to \
2615perform, and such a rule changes nothing even once applied. \
2616(2) {\"kind\":\"fact\",\"relation\":\"...\",\"object\":\"...\"} with the same \
2617entity target — a durable fact the agent keeps having to be told (an alias, a \
2618settled default, a preference). 'relation' is a short identifier (letters, \
2619digits, _ - . :), not a sentence. \
2620(3) {\"kind\":\"query_revision\",\"body\":\"<CAL>\"} with target \
2621\"query:<name>\" or \"template:<name>\" — a rewrite of the saved query that \
2622assembles the agent's context, when the evidence shows it retrieves the wrong \
2623things. Give the FULL new body; it replaces the old one. \
2624(4) {\"kind\":\"plan_revision\",\"edits\":[{\"path\":\"...\",\"from\":X,\"to\":Y}]} \
2625with target \"grain:<workflow hash>\" — at most 8 field-level edits to the \
2626workflow. Only these paths are editable: 'edges.<i>.cond', \
2627'edges.<i>.max_cycles', 'retries.<node>'. 'from' MUST equal what the plan \
2628holds now, or the edit is refused. You cannot add, remove or rewire nodes. \
2629(5) {\"kind\":\"code_revision\",\"source\":\"...\"} with target \"tool:<name>\" \
2630— full replacement source for that tool. It is applied only after a recorded \
2631evaluation run passes, so propose one only when the evidence shows the current \
2632code is the defect. \
2633The subject of a fact, the name of a query, the plan hash and the tool name \
2634all come from 'target' — do not repeat them inside 'proposal'. A proposal \
2635becomes a change a human reviewer may apply, so it must be fully supported by \
2636the cited evidence. Propose nothing you cannot ground in the evidence."
2637        )
2638    };
2639}
2640
2641/// The review-queue objective (the default).
2642const DISCOVER_INSTRUCTIONS: &str = discover_instructions!(
2643    "SCORING: propose a finding ONLY if you \
2644are more than 0.75 confident it is BOTH correct AND materially useful. A correct, \
2645useful finding earns 1; a wrong or trivial one is penalized 2; returning nothing \
2646earns 0. When in doubt, propose nothing — an empty list is the correct answer \
2647when there is nothing worth flagging."
2648);
2649
2650/// The learner objective (`Policy::discover_objective = learner`).
2651const DISCOVER_LEARNER_INSTRUCTIONS: &str = discover_instructions!(
2652    "SCORING: you are the learning stage of a deployed agent, and what you \
2653propose now is what it will do differently next time — a lesson you withhold \
2654is a mistake it repeats. A correct, actionable proposal earns 1; a wrong or \
2655trivial one is penalized 1; returning nothing while the evidence holds a \
2656recurring failure, two or more rejected outcomes, an instruction from a \
2657person, or a multi-step procedure the agent completed successfully that no \
2658saved skill or plan covers, is ALSO penalized 1. Abstain only when the \
2659evidence shows none of those. Prefer the one proposal that addresses the most \
2660frequent or most costly failure — or, when nothing failed, the procedure that \
2661worked — over several speculative ones, and report your confidence honestly — \
2662an independent verifier, not you, decides what survives."
2663);
2664
2665/// The fixed GROUND instruction (§5.2): verify the finding's factual PREMISES
2666/// are real (anti-fabrication), while allowing an inference. A self-improvement
2667/// finding may reason BEYOND the evidence (e.g. 'HQ=SF and country=Germany are
2668/// inconsistent'); grounding checks the premises (HQ=SF, country=Germany) are in
2669/// the evidence — the soundness of the inference is VERIFY's job, not this one.
2670const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2671fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2672the facts it relies on are actually present in the cited evidence, NOT that its \
2673conclusion is stated verbatim. Decompose the finding into the factual claims it \
2674depends on. Mark supported=true when those facts are present in the evidence \
2675(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2676on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2677different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2678
2679/// The fixed VERIFY instruction (§5.3): adversarial, abstention-biased.
2680const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2681each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2682never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2683SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2684'possible' findings with no concrete defect, and reject any claimed \
2685inconsistency or contradiction that is not backed by at least two actually \
2686conflicting facts in the cited evidence. (2) Context — does the finding \
2687correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2688finding when it names a genuine, specific problem grounded in its evidence and \
2689materially useful to a human reviewer; otherwise reject it, and default to \
2690keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2691redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2692grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2693KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2694{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2695
2696/// The fixed ENRICH instruction.
2697const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2698guidance note to help a human reviewer decide. Do not restate the finding. Return \
2699JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2700
2701/// Add a grain to the evidence bundle — deduplicated, bounded at 64, text
2702/// capped. Shared by the deterministic-citation and recent-grain seeding.
2703/// `ns_by_hash` records each bundled grain's namespace so an authored lesson
2704/// can later land in the namespace its evidence lives in (never one the
2705/// model names).
2706fn push_evidence(
2707    evidence: &mut Vec<crate::llm::EvidenceItem>,
2708    bundle: &mut BTreeSet<String>,
2709    ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2710    g: &GrainRecord,
2711    attribution: crate::policy::EvidenceAttribution,
2712) {
2713    if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2714        ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2715        evidence.push(crate::llm::EvidenceItem {
2716            id: format!("e{}", evidence.len() + 1),
2717            hash: g.hash.clone(),
2718            grain_type: g.grain_type.clone(),
2719            text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
2720        });
2721    }
2722}
2723
2724/// Resolve one citation to a bundled grain's hash: the full hash, the
2725/// bundle-local `id` the evidence item carried, or an unambiguous hash prefix
2726/// of at least 12 hex chars. Anything else is a fabrication and resolves to
2727/// nothing. Small models copy 64-hex hashes badly — measured live, the
2728/// cite-check was where most of a cheap model's drafts died ("proposed 3 →
2729/// cited 1") — and a citation that plainly names one bundled grain is not
2730/// the thing that check exists to catch.
2731pub(crate) fn resolve_citation(
2732    cite: &str,
2733    bundle: &BTreeSet<String>,
2734    id_to_hash: &std::collections::BTreeMap<&str, &str>,
2735) -> Option<String> {
2736    let cite = cite.trim();
2737    if bundle.contains(cite) {
2738        return Some(cite.to_string());
2739    }
2740    if let Some(h) = id_to_hash.get(cite) {
2741        return Some((*h).to_string());
2742    }
2743    const MIN_PREFIX: usize = 12;
2744    if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
2745        let lower = cite.to_ascii_lowercase();
2746        let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
2747        if let (Some(h), None) = (it.next(), it.next()) {
2748            return Some(h.clone());
2749        }
2750    }
2751    None
2752}
2753
2754/// A short human-readable projection of a grain for the evidence bundle,
2755/// under the host's attribution policy.
2756fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
2757    if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2758        return format!("{s} {r} {o}");
2759    }
2760    // Tool grains — the evidence most lesson drafts cite. Rendering them
2761    // empty starved GROUND of the very facts it exists to check: a correct
2762    // lesson would be refused as unverifiable (found live — the gate
2763    // rightly rejected a claim over evidence it could not see).
2764    if let Some(t) = g.tool_name() {
2765        let status = if g.is_error() { "error" } else { "ok" };
2766        let out = g.tool_content().unwrap_or("");
2767        // The call's input, when recorded: without it a successful trajectory
2768        // reads as a list of tool names and outputs, and a procedure — WHICH
2769        // ticket was fetched, WHAT tag was set — cannot be reconstructed from
2770        // it. A skill proposal needs the arguments; a lesson usually does not,
2771        // and the cap on the brief bounds the cost either way.
2772        let input = match g.fields.get("input") {
2773            Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
2774            Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
2775            _ => String::new(),
2776        };
2777        return format!("tool {t}{input} {status}: {out}");
2778    }
2779    // `object` is last but it is not optional: an Observation stores its text
2780    // there (subject + object, no relation), so it misses the fact-triple
2781    // branch above and used to fall through this list to an empty string —
2782    // every human note in a memory reached the model as a blank line. That is
2783    // the single highest-value evidence a memory holds, and it was the one
2784    // shape that rendered to nothing. Callers had started duplicating the text
2785    // into `body` to work around it; nothing should have to.
2786    for key in ["content", "body", "text", "summary", "object"] {
2787        if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2788            if v.is_empty() {
2789                continue;
2790            }
2791            // Who said it, when the grain records it. An Observation reaches
2792            // the model as a bare sentence otherwise, and a bare sentence is
2793            // ambiguous about direction in exactly the way that matters: a
2794            // person's correction ("Vendor Name is ACME") reads identically
2795            // to the agent having been told something it asked for. Measured
2796            // live on the receipts corpus, a model given 31 unattributed
2797            // corrections concluded the agent was repeatedly *requesting*
2798            // data it already had, and proposed rules to stop it asking.
2799            // The observer is already on the grain; only the projection
2800            // dropped it.
2801            if attribution == crate::policy::EvidenceAttribution::Anonymous {
2802                return v.to_string();
2803            }
2804            if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
2805                if !who.is_empty() {
2806                    let kind = g
2807                        .fields
2808                        .get("observer_type")
2809                        .and_then(|v| v.as_str())
2810                        .unwrap_or("");
2811                    let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
2812                    let mut prefix = if kind == "human" {
2813                        format!("{who} (a person) said")
2814                    } else {
2815                        format!("{who} observed")
2816                    };
2817                    if !about.is_empty() {
2818                        prefix.push_str(&format!(" of {about}"));
2819                    }
2820                    return format!("{prefix}: {v}");
2821                }
2822            }
2823            return v.to_string();
2824        }
2825    }
2826    String::new()
2827}
2828
2829/// One imperative line, no control characters, capped: the only shape an
2830/// authored lesson may take. A literal newline could otherwise smuggle a
2831/// second CAL statement past review (belt: serde_json escapes it anyway) or
2832/// break the one-line prompt rendering hosts assume.
2833fn sanitize_lesson(s: &str) -> String {
2834    sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2835}
2836
2837/// One line, no control characters, capped. Every free-text field the model
2838/// can put into an executable proposal goes through this: a literal newline
2839/// could otherwise smuggle a second CAL statement past review (belt:
2840/// serde_json escapes it anyway), split a one-line DEFINE across the batch the
2841/// apply path iterates, or break the one-line prompt rendering hosts assume.
2842fn sanitize_line(s: &str, max: usize) -> String {
2843    let cleaned: String =
2844        s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2845    crate::llm::cap(cleaned.trim(), max)
2846}
2847
2848/// A relation is an identifier, not prose — it becomes a queryable predicate,
2849/// and whitespace or quotes in one would make the Fact unfindable by the very
2850/// recall that should surface it. `None` rejects the draft's `fact` proposal.
2851fn sanitize_relation(s: &str) -> Option<String> {
2852    let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2853    if r.is_empty()
2854        || !r
2855            .chars()
2856            .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2857    {
2858        return None;
2859    }
2860    Some(r)
2861}
2862
2863/// A model-supplied query body is placed INSIDE a `DEFINE … AS { … }` block,
2864/// so it is the one place in the vocabulary where model text becomes part of a
2865/// statement's structure rather than its content. Two belts, because the
2866/// substrate's parser strength is not something this engine gets to assume:
2867///
2868/// - **No braces.** Closing the block early is the injection shape; a saved
2869///   RECALL/ASSEMBLE body needs no braces of its own, so refusing them costs
2870///   nothing and fails closed.
2871/// - **No destructive keyword, anywhere in the body.** `cal::contains_destructive`
2872///   scans each LINE's leading keyword, which a single-line injection slips
2873///   past by construction — so scan every token here instead.
2874///
2875/// The substrate's own `validate_cal` and the saved-query read-only
2876/// verification pass still run after this; this is the layer that does not
2877/// depend on either of them being strict.
2878fn safe_definition_body(body: &str) -> bool {
2879    if body.contains('{') || body.contains('}') {
2880        return false;
2881    }
2882    !body
2883        .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2884        .any(|tok| {
2885            ["FORGET", "PURGE", "DROP", "DEFINE"]
2886                .iter()
2887                .any(|kw| tok.eq_ignore_ascii_case(kw))
2888        })
2889}
2890
2891/// The claim GROUND entails and VERIFY stress-tests. When the draft carries a
2892/// resolvable proposal the claim names exactly what an apply would do, so what
2893/// survives the gates is what gets written — never a summary standing in for
2894/// a change the gates never saw.
2895fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2896    let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2897    match resolved {
2898        Some(r) => format!("{summary} {}", r.rendered),
2899        None => summary,
2900    }
2901}
2902
2903/// A DISCOVER draft that survived structural validation, carrying the
2904/// executable form of its proposal. Resolution happens BEFORE GROUND/VERIFY,
2905/// so both gates judge exactly what an apply would do — the rule the authored
2906/// lesson already followed, generalized to the whole vocabulary. It also means
2907/// a malformed proposal costs no model call: it dies here, not at apply.
2908struct ValidatedDraft {
2909    draft: crate::llm::LlmDraft,
2910    target_ref: String,
2911    cited: Vec<String>,
2912    resolved: Option<ResolvedProposal>,
2913}
2914
2915/// The executable shape of a validated draft. `None` on a [`ValidatedDraft`]
2916/// means advisory — the model said something a human may want to see, but
2917/// nothing the engine will ever execute.
2918struct ResolvedProposal {
2919    action: ActionKind,
2920    proposal: Proposal,
2921    /// One line naming exactly what an apply would do; folded into the claim
2922    /// both gates judge and shown in the review summary.
2923    rendered: String,
2924    summary_key: &'static str,
2925    summary_args: serde_json::Map<String, Value>,
2926    rollbackable: bool,
2927    evalset_hash: Option<String>,
2928    importance: f64,
2929    /// Set for the two Fact-writing shapes. Held rather than pre-rendered so
2930    /// the grain can carry the VERIFIER's confidence, which is not known until
2931    /// after the gates have run.
2932    fact_fields: Option<serde_json::Map<String, Value>>,
2933    /// Statements appended after the `fact_fields` ADD when the proposal is
2934    /// rendered — a consolidation's supersessions of the pile it replaces.
2935    extra_statements: Vec<String>,
2936    /// A plan revision's rehearsal report (`SubstrateRead::plan_replay`),
2937    /// carried onto the recommendation so the reviewer sees it.
2938    replay: Option<Value>,
2939}
2940
2941/// The Workflow fields a `plan_revision` may touch. Thresholds and limits —
2942/// never who calls what.
2943///
2944/// The exclusion is structural, not advisory: `nodes`, `edges[].src`,
2945/// `edges[].dst` and `bindings` are simply not matchable here, so a topology
2946/// change cannot be expressed by any proposal the model can write. That keeps
2947/// a plan revision reviewable as a short list of scalar deltas rather than a
2948/// re-drawn graph, which is the difference between a reviewer checking a
2949/// number and a reviewer re-deriving a plan.
2950fn plan_edit_allowed(path: &str) -> bool {
2951    let seg: Vec<&str> = path.split('.').collect();
2952    match seg.as_slice() {
2953        ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2954        ["retries", node] => !node.is_empty(),
2955        _ => false,
2956    }
2957}
2958
2959/// Read the value at an allowlisted path (absent → `Value::Null`, which is
2960/// what an edit adding a `retries` entry must declare as its `from`).
2961fn plan_get(body: &Value, path: &str) -> Value {
2962    let mut cur = body;
2963    for seg in path.split('.') {
2964        cur = match cur {
2965            Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2966                Some(v) => v,
2967                None => return Value::Null,
2968            },
2969            Value::Object(o) => match o.get(seg) {
2970                Some(v) => v,
2971                None => return Value::Null,
2972            },
2973            _ => return Value::Null,
2974        };
2975    }
2976    cur.clone()
2977}
2978
2979/// Write the value at an allowlisted path. Only creates a missing key in an
2980/// object (the `retries.<node>` case — including the `retries` map itself,
2981/// which a stored plan with no retries omits entirely under omit-defaults
2982/// serialization, so a first retry edit on such a plan used to fail to
2983/// resolve and stay advisory); never grows an array.
2984fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2985    let segs: Vec<&str> = path.split('.').collect();
2986    let Some((last, parents)) = segs.split_last() else {
2987        return false;
2988    };
2989    let mut cur = body;
2990    for (depth, seg) in parents.iter().enumerate() {
2991        cur = match cur {
2992            Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2993                Some(v) => v,
2994                None => return false,
2995            },
2996            Value::Object(o) => {
2997                if depth == 0 && *seg == "retries" && !o.contains_key("retries") {
2998                    o.insert("retries".into(), Value::Object(serde_json::Map::new()));
2999                }
3000                match o.get_mut(*seg) {
3001                    Some(v) => v,
3002                    None => return false,
3003                }
3004            }
3005            _ => return false,
3006        };
3007    }
3008    match cur {
3009        Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
3010            Some(slot) => {
3011                *slot = to;
3012                true
3013            }
3014            None => false,
3015        },
3016        Value::Object(o) => {
3017            o.insert((*last).to_string(), to);
3018            true
3019        }
3020        _ => false,
3021    }
3022}
3023
3024/// Type-check one plan edit's new value against the field it targets. Without
3025/// this a string in `max_cycles` would be dropped by the grain deserializer
3026/// and the "applied" revision would silently mean *unlimited* — a proposal
3027/// that reads as a tightening and lands as a removal.
3028fn plan_value_ok(path: &str, to: &Value) -> bool {
3029    let seg: Vec<&str> = path.split('.').collect();
3030    match seg.as_slice() {
3031        ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
3032            !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
3033        }),
3034        ["edges", _, "max_cycles"] | ["retries", _] => {
3035            to.as_u64().is_some_and(|n| n <= 1_000)
3036        }
3037        _ => false,
3038    }
3039}
3040
3041/// Resolve a DISCOVER draft's proposal into the executable form an apply would
3042/// run, or `None` for advisory. Every variant takes its SCOPE from the draft's
3043/// target and its supporting facts from the substrate — the model names the
3044/// change, never the subject it lands on, the namespace it lands in, or (for
3045/// code) the evalset that grades it.
3046fn resolve_proposal<S: OmsSubstrate>(
3047    sub: &S,
3048    d: &crate::llm::LlmDraft,
3049    target: &TargetRef,
3050    cited: &[String],
3051    ns_by_hash: &std::collections::BTreeMap<String, String>,
3052    caps: Capabilities,
3053    policy: &crate::policy::Policy,
3054) -> Option<ResolvedProposal> {
3055    use crate::llm::DraftProposal as P;
3056    let (skills, plans) = (&policy.skills, &policy.plans);
3057    let mut args = serde_json::Map::new();
3058    match d.parsed_proposal()? {
3059        // ---- plan: a procedure as a validated Workflow + its Skill prose ----
3060        P::Plan { description, when_to_use, nodes, edges } => {
3061            if !plans.enabled || !caps.plans {
3062                return None;
3063            }
3064            let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
3065                derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
3066            args.insert("name".into(), Value::from(name.clone()));
3067            args.insert("nodes".into(), Value::from(n_nodes as u64));
3068            let mut stmts = vec![match &existing_skill {
3069                Some(h) => cal::supersede(h, "skill", &skill),
3070                None => cal::add("skill", &skill),
3071            }];
3072            // A graph the runtime would refuse is not minted as a plan — but
3073            // the procedure it describes is still the thing worth keeping, so
3074            // it is recorded as a skill. Measured need: a 30B proposer writes
3075            // conditions like "tickets.length > 0", outside the frozen v1
3076            // grammar, and discarding the draft for that threw away the
3077            // captured procedure entirely (PERSIST.md §11 #27).
3078            let (summary_key, kind) = match &workflow {
3079                Some(wf) => {
3080                    stmts.push(match &existing_plan {
3081                        Some(h) => cal::supersede(h, "workflow", wf),
3082                        None => cal::add("workflow", wf),
3083                    });
3084                    args.insert("edges".into(), Value::from(n_edges as u64));
3085                    ("llm.plan", "plan")
3086                }
3087                None => {
3088                    args.insert("steps".into(), Value::from(n_nodes as u64));
3089                    ("llm.skill", "skill")
3090                }
3091            };
3092            let patched = existing_skill.is_some() || (workflow.is_some() && existing_plan.is_some());
3093            let (action, verb) = if patched {
3094                (ActionKind::Revise, "revise")
3095            } else {
3096                (ActionKind::Record, "record")
3097            };
3098            Some(ResolvedProposal {
3099                action,
3100                proposal: Proposal::Cal { cal: cal::batch(&stmts) },
3101                rendered: format!(
3102                    "Proposed {kind} to {verb}: \"{name}\" — {n_nodes} steps; when: {}",
3103                    skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
3104                ),
3105                summary_key,
3106                summary_args: args,
3107                rollbackable: true,
3108                evalset_hash: None,
3109                importance: 0.65,
3110                fact_fields: None,
3111                extra_statements: Vec::new(),
3112                replay: None,
3113            })
3114        }
3115        // ---- skill: a reusable procedure from a trajectory that succeeded ----
3116        P::Skill { description, when_to_use, steps } => {
3117            if !skills.enabled {
3118                return None;
3119            }
3120            let SkillFields { fields, name, n_steps, existing } =
3121                derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
3122            args.insert("name".into(), Value::from(name.clone()));
3123            args.insert("steps".into(), Value::from(n_steps as u64));
3124            // A live skill of the same name in the same namespace is PATCHED
3125            // (superseded), never duplicated beside itself.
3126            let (action, cal, verb) = match existing {
3127                Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
3128                None => (ActionKind::Record, cal::add("skill", &fields), "record"),
3129            };
3130            Some(ResolvedProposal {
3131                action,
3132                proposal: Proposal::Cal { cal },
3133                rendered: format!(
3134                    "Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
3135                    fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
3136                ),
3137                summary_key: "llm.skill",
3138                summary_args: args,
3139                rollbackable: true,
3140                evalset_hash: None,
3141                importance: 0.6,
3142                fact_fields: None,
3143                extra_statements: Vec::new(),
3144                replay: None,
3145            })
3146        }
3147        // ---- consolidation: one lesson replacing a pile ----
3148        P::Consolidation { lesson, supersedes } => {
3149            let lesson = sanitize_lesson(&lesson);
3150            if lesson.is_empty() {
3151                return None;
3152            }
3153            let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
3154            let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("").to_string();
3155            // Every member must be a LIVE lesson on this very entity, all in
3156            // one namespace, and there must be a pile — a "consolidation" of
3157            // one grain, or of grains the model picked from elsewhere, is
3158            // not a consolidation and stays advisory.
3159            let mut hashes: Vec<String> = supersedes.into_iter().collect();
3160            hashes.sort();
3161            hashes.dedup();
3162            let mut members = Vec::new();
3163            for h in &hashes {
3164                let g = sub.grain(h).ok().flatten()?;
3165                if !g.is_live()
3166                    || g.fact_relation() != Some("lesson")
3167                    || g.fact_subject().is_none_or(|s| normalize_ident(s) != normalize_ident(&subject))
3168                {
3169                    return None;
3170                }
3171                members.push(g);
3172            }
3173            if members.len() < 2 || members.iter().any(|m| m.namespace != members[0].namespace) {
3174                return None;
3175            }
3176            let ns = members[0].namespace.clone();
3177            let mut fields = fields;
3178            if !ns.is_empty() {
3179                fields.insert("namespace".into(), Value::from(ns.clone()));
3180            }
3181            fields.insert("consolidates".into(), Value::from(hashes.clone()));
3182            // Each member is superseded by a marker naming the line that
3183            // replaced it — not by a copy of the lesson, which would leave N
3184            // live copies in the prompt. Rollback retracts the markers and
3185            // the added line; the members come back as heads.
3186            let extra_statements: Vec<String> = members
3187                .iter()
3188                .map(|m| {
3189                    let mut marker = serde_json::Map::new();
3190                    marker.insert("subject".into(), Value::from(subject.clone()));
3191                    marker.insert("relation".into(), Value::from("mg:lesson_consolidated"));
3192                    marker.insert("object".into(), Value::from(lesson.clone()));
3193                    if !ns.is_empty() {
3194                        marker.insert("namespace".into(), Value::from(ns.clone()));
3195                    }
3196                    cal::supersede(&m.hash, "fact", &marker)
3197                })
3198                .collect();
3199            args.insert("lesson".into(), Value::from(lesson.clone()));
3200            args.insert("count".into(), Value::from(members.len() as u64));
3201            Some(ResolvedProposal {
3202                action: ActionKind::Consolidate,
3203                proposal: Proposal::Cal { cal: cal::batch(&extra_statements) },
3204                rendered: format!(
3205                    "Proposed consolidation of {} lessons on \"{subject}\" into one: \"{lesson}\"",
3206                    members.len()
3207                ),
3208                summary_key: "llm.consolidation",
3209                summary_args: args,
3210                rollbackable: true,
3211                evalset_hash: None,
3212                importance: 0.6,
3213                fact_fields: Some(fields),
3214                extra_statements,
3215                replay: None,
3216            })
3217        }
3218        // ---- lesson: the pre-vocabulary shape, unchanged ----
3219        P::Lesson { lesson } => {
3220            let lesson = sanitize_lesson(&lesson);
3221            if lesson.is_empty() {
3222                return None;
3223            }
3224            let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
3225            args.insert("lesson".into(), Value::from(lesson.clone()));
3226            Some(ResolvedProposal {
3227                // Same action as the deterministic lesson path — "record a
3228                // failure-derived lesson" — so dedup groups authored lessons
3229                // per target and review UIs need no new vocabulary.
3230                action: ActionKind::ClusterFailure,
3231                proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
3232                rendered: format!("Proposed lesson to record: \"{lesson}\""),
3233                summary_key: "llm.lesson",
3234                summary_args: args,
3235                rollbackable: true,
3236                evalset_hash: None,
3237                importance: 0.5,
3238                fact_fields: Some(fields),
3239                extra_statements: Vec::new(),
3240                replay: None,
3241            })
3242        }
3243        // ---- fact: a durable fact under a model-chosen relation ----
3244        P::Fact { relation, object } => {
3245            let relation = sanitize_relation(&relation)?;
3246            let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
3247            if object.is_empty() {
3248                return None;
3249            }
3250            let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
3251            let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
3252            args.insert("relation".into(), Value::from(relation.clone()));
3253            args.insert("object".into(), Value::from(object.clone()));
3254            Some(ResolvedProposal {
3255                action: ActionKind::Record,
3256                proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
3257                rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
3258                summary_key: "llm.fact",
3259                summary_args: args,
3260                rollbackable: true,
3261                evalset_hash: None,
3262                importance: 0.5,
3263                fact_fields: Some(fields),
3264                extra_statements: Vec::new(),
3265                replay: None,
3266            })
3267        }
3268        // ---- query_revision: how the agent assembles its own context ----
3269        P::QueryRevision { body } => {
3270            let name = target.opaque();
3271            // The name comes from the target, but it still ends up inside a
3272            // statement — a quote or control character in one would change the
3273            // statement's shape rather than its content.
3274            if name.is_empty()
3275                || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
3276            {
3277                return None;
3278            }
3279            let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
3280            if body.is_empty() || !safe_definition_body(&body) {
3281                return None;
3282            }
3283            let stmt = match target.scheme() {
3284                "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
3285                "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
3286                _ => return None,
3287            };
3288            // The substrate owns the grammar: if it will not parse, or will
3289            // not hand back an inverse, this is not something a reviewer
3290            // should be offered as applicable. A definition change ROLLBACK
3291            // could not undo must not be applied at all.
3292            sub.validate_cal(&stmt).ok()?;
3293            sub.definition_inverse(&stmt).ok().flatten()?;
3294            args.insert("name".into(), Value::from(name));
3295            args.insert("body".into(), Value::from(body.clone()));
3296            Some(ResolvedProposal {
3297                action: ActionKind::Revise,
3298                proposal: Proposal::Cal { cal: stmt },
3299                rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
3300                summary_key: "llm.query_revision",
3301                summary_args: args,
3302                rollbackable: true,
3303                evalset_hash: None,
3304                importance: 0.6,
3305                fact_fields: None,
3306                extra_statements: Vec::new(),
3307                replay: None,
3308            })
3309        }
3310        // ---- plan_revision: field-level edits to a Workflow grain ----
3311        P::PlanRevision { edits } => {
3312            if !caps.plans
3313                || target.scheme() != "grain"
3314                || edits.is_empty()
3315                || edits.len() > crate::llm::MAX_PLAN_EDITS
3316            {
3317                return None;
3318            }
3319            let hash = target.opaque();
3320            let g = sub.grain(hash).ok().flatten()?;
3321            if g.grain_type != "workflow" || !g.is_live() {
3322                return None;
3323            }
3324            let mut body = Value::Object(g.fields.clone());
3325            let mut deltas = Vec::new();
3326            let nodes: std::collections::BTreeSet<String> = body
3327                .get("nodes")
3328                .and_then(Value::as_array)
3329                .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
3330                .unwrap_or_default();
3331            for e in &edits {
3332                if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
3333                    return None;
3334                }
3335                // A retry count for a node that does not exist is inert, but
3336                // applying it still mints a new plan hash — and every trigger
3337                // pointing at the old one must then be walked forward. A
3338                // no-op is not worth that.
3339                if let Some(node) = e.path.strip_prefix("retries.") {
3340                    if !nodes.contains(node) {
3341                        return None;
3342                    }
3343                }
3344                // `from` is the staleness check: a proposal authored against
3345                // an older plan does not silently apply to a newer one.
3346                if plan_get(&body, &e.path) != e.from {
3347                    return None;
3348                }
3349                // A no-op edit is not a revision; it would apply, mint a new
3350                // plan hash, and orphan every trigger pointing at the old one
3351                // for nothing.
3352                if e.from == e.to {
3353                    return None;
3354                }
3355                if !plan_set(&mut body, &e.path, e.to.clone()) {
3356                    return None;
3357                }
3358                deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
3359            }
3360            // The substrate owns the plan grammar (unique + reachable nodes,
3361            // conditions parse, every cycle bounded). An edit that would make
3362            // the plan unrunnable never reaches a reviewer as applicable.
3363            sub.validate_plan(&body).ok()?;
3364            // The rehearsal: the candidate re-driven through the runtime's
3365            // scheduler over the live plan's journaled runs, every effect
3366            // answered from the journal (`areev run shadow --plan-file`).
3367            // The report rides on the card either way; under a
3368            // `plan_replay` policy it is also the gate — a candidate worse
3369            // than the incumbent on the same runs, or rehearsable on too
3370            // few of them, is stored as advisory with the reason naming the
3371            // runs, never offered to apply.
3372            let replay = sub.plan_replay(hash, &body).ok().flatten();
3373            let refused = match (&policy.plan_replay, &replay) {
3374                (Some(gate), Some(report)) => gate.refusal(report),
3375                _ => None,
3376            };
3377            let Value::Object(fields) = body else {
3378                return None;
3379            };
3380            args.insert("plan".into(), Value::from(hash));
3381            args.insert("edits".into(), Value::from(deltas.join("; ")));
3382            if let Some(reason) = refused {
3383                let mut data = serde_json::Map::new();
3384                data.insert("plan".into(), Value::from(hash));
3385                data.insert("edits".into(), Value::from(deltas.clone()));
3386                data.insert("refused".into(), Value::from(reason.clone()));
3387                args.insert("reason".into(), Value::from(reason.clone()));
3388                return Some(ResolvedProposal {
3389                    action: ActionKind::Flag,
3390                    proposal: Proposal::Data { data },
3391                    rendered: format!(
3392                        "Plan revision ({}) refused by the rehearsal: {reason}",
3393                        deltas.join("; ")
3394                    ),
3395                    summary_key: "llm.plan_revision_refused",
3396                    summary_args: args,
3397                    rollbackable: false,
3398                    evalset_hash: None,
3399                    importance: 0.4,
3400                    fact_fields: None,
3401                    extra_statements: Vec::new(),
3402                    replay,
3403                });
3404            }
3405            let stmt = cal::supersede(hash, "workflow", &fields);
3406            // Same rule as the definition rewrite: a statement the substrate
3407            // will not accept is not something to offer a reviewer as
3408            // applicable. `validate_plan` checked the GRAPH; this checks the
3409            // statement that carries it.
3410            sub.validate_cal(&stmt).ok()?;
3411            Some(ResolvedProposal {
3412                action: ActionKind::Revise,
3413                proposal: Proposal::Cal { cal: stmt },
3414                rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
3415                summary_key: "llm.plan_revision",
3416                summary_args: args,
3417                rollbackable: true,
3418                evalset_hash: None,
3419                importance: 0.7,
3420                fact_fields: None,
3421                extra_statements: Vec::new(),
3422                replay,
3423            })
3424        }
3425        // ---- code_revision: §7.4, gated by the tool's own evalset ----
3426        P::CodeRevision { source } => {
3427            if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
3428                return None;
3429            }
3430            if source.chars().count() > crate::llm::MAX_CODE_LEN {
3431                return None;
3432            }
3433            // Rule E1's pin, resolved from the substrate. A proposer that
3434            // could name its own grader is not gated, and a tool that
3435            // declares no evalset has no gate to pass — advisory either way.
3436            let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
3437            let mut data = serde_json::Map::new();
3438            data.insert("tool".into(), Value::from(target.opaque()));
3439            data.insert("source".into(), Value::from(source.clone()));
3440            args.insert("tool".into(), Value::from(target.opaque()));
3441            args.insert("bytes".into(), Value::from(source.len() as u64));
3442            Some(ResolvedProposal {
3443                action: ActionKind::CodeRevision,
3444                proposal: Proposal::Data { data },
3445                rendered: format!(
3446                    "Proposed new source for tool {} ({} bytes), gated by evalset {}",
3447                    target.opaque(),
3448                    source.len(),
3449                    evalset
3450                ),
3451                summary_key: "llm.code_revision",
3452                summary_args: args,
3453                rollbackable: true,
3454                evalset_hash: Some(evalset),
3455                importance: 0.8,
3456                fact_fields: None,
3457                extra_statements: Vec::new(),
3458                replay: None,
3459            })
3460        }
3461    }
3462}
3463
3464/// Stamp a validated DISCOVER draft as an `origin = llm` recommendation.
3465/// Default shape: an advisory `Flag` carrying `Proposal::Data` (no executable
3466/// mutation). A draft whose proposal RESOLVED (see [`resolve_proposal`])
3467/// instead stamps as that executable change, reviewable with the exact line an
3468/// apply would run. Either way `Origin::Llm` plus the no-manifest analyzer id
3469/// leave it structurally ineligible for auto-apply — and independently, no
3470/// class this vocabulary can reach except `memory` is auto-appliable at all
3471/// (`Policy::grants_auto_apply`). The only path into the agent is a human
3472/// review with a BECAUSE followed by an explicit apply.
3473fn stamp_llm(
3474    model: &str,
3475    d: &crate::llm::LlmDraft,
3476    target_ref: String,
3477    cited: Vec<String>,
3478    resolved: Option<ResolvedProposal>,
3479    confidence: f64,
3480    now_ms: i64,
3481) -> Recommendation {
3482    let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3483    let guidance = if d.guidance.trim().is_empty() {
3484        None
3485    } else {
3486        Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
3487    };
3488    let replay = resolved.as_ref().and_then(|r| r.replay.clone());
3489    let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
3490        Some(mut r) => {
3491            // What the proposal would DO, before the verifier's confidence
3492            // is folded into the fact: the dedup key fingerprints this, so
3493            // the same lesson at a different confidence is one finding.
3494            let content = match &r.fact_fields {
3495                Some(fields) => format!(
3496                    "{} {}",
3497                    fields.get("relation").and_then(Value::as_str).unwrap_or(""),
3498                    fields.get("object").and_then(Value::as_str).unwrap_or("")
3499                ),
3500                None => match &r.proposal {
3501                    Proposal::Cal { cal } => cal.clone(),
3502                    Proposal::Data { data } => Value::Object(data.clone()).to_string(),
3503                    Proposal::Edit { diff, .. } => diff.clone(),
3504                },
3505            };
3506            // The grain records the VERIFIER's calibrated confidence — the
3507            // independent signal — never the proposer's self-report.
3508            if let Some(mut fields) = r.fact_fields.take() {
3509                fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
3510                let mut statements = vec![cal::add("fact", &fields)];
3511                statements.extend(r.extra_statements.iter().cloned());
3512                r.proposal = Proposal::Cal { cal: cal::batch(&statements) };
3513            }
3514            let mut args = r.summary_args;
3515            args.insert("text".into(), Value::from(summary_text));
3516            (
3517                r.action,
3518                r.proposal,
3519                Summary::new(r.summary_key, args),
3520                r.rollbackable,
3521                r.importance,
3522                r.evalset_hash,
3523                Some(content),
3524            )
3525        }
3526        None => {
3527            let mut args = serde_json::Map::new();
3528            args.insert("text".into(), Value::from(summary_text));
3529            let mut data = serde_json::Map::new();
3530            data.insert("source".into(), Value::from("llm"));
3531            (
3532                ActionKind::Flag,
3533                Proposal::Data { data },
3534                Summary::new("llm.discover", args),
3535                false,
3536                0.3,
3537                None,
3538                None,
3539            )
3540        }
3541    };
3542    // An advisory flag keeps the analyzer-style key (one open flag per
3543    // target); an executable proposal keys on its content too, because
3544    // there the content is the finding.
3545    let dedup = match &content {
3546        Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
3547        None => dedup_key("llm", &target_ref, action),
3548    };
3549    Recommendation {
3550        hash: String::new(),
3551        analyzer: "loop.llm/1".to_string(),
3552        params_snapshot: serde_json::Map::new(),
3553        origin: Origin::Llm { model: model.to_string() },
3554        target_ref: target_ref.clone(),
3555        action_kind: action,
3556        dedup_key: dedup,
3557        summary,
3558        severity: Severity::Low,
3559        proposal,
3560        destructive: false,
3561        rollbackable,
3562        evidence: cited,
3563        evidence_query: None,
3564        metric: None,
3565        // The verifier's calibrated confidence — not a hardcoded default.
3566        confidence: confidence.clamp(0.0, 1.0),
3567        importance,
3568        created_at_ms: now_ms,
3569        guidance,
3570        evalset_hash,
3571        near_duplicate_of: Vec::new(),
3572        replay,
3573        status: RecStatus::Pending,
3574    }
3575}
3576
3577/// The Fact an approved `lesson` or `fact` proposal writes: subject from the
3578/// entity target, the model's relation (`"lesson"` for the lesson shape —
3579/// prescriptive prose, distinct from the deterministic `fails_with` signature
3580/// facts), the sanitized object, and the DOMINANT namespace of the cited
3581/// evidence (max count, ties to the lexicographically smallest — the
3582/// tool_failure rule), never a namespace the model names. `None` when the
3583/// target gives no subject.
3584/// The `skill` paragraph appended to the DISCOVER instructions when the host
3585/// allows skill authoring. Kept beside the fixed instruction text it extends.
3586fn skill_instructions(min_steps: u32) -> String {
3587    format!(
3588        " (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
3589\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
3590PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
3591calls that reached its goal, which a later session facing the same situation \
3592should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
3593the tool called and the values that mattered (the field checked, the tag set, the \
3594exact format produced), a one-line description, and 'when_to_use' — the situation \
3595that should trigger it. The skill-name is a short identifier (letters, digits, \
3596_ -). If a saved skill already covers this procedure, use ITS name so it is \
3597patched rather than duplicated. Do not propose a skill for a procedure that \
3598failed, or for one already saved and unchanged. A finding that itself describes \
3599two or more steps the agent should carry out in order ('after listing the \
3600tickets, fetch each, then …') IS a procedure: propose it as a skill or a plan, \
3601never as a lesson — a lesson is one rule, and a procedure written as one is a \
3602procedure nobody can open.",
3603        crate::llm::MAX_SKILL_STEPS
3604    )
3605}
3606
3607/// The `plan` paragraph appended to the DISCOVER instructions when the host
3608/// allows plan authoring. The condition grammar is the runtime's frozen v1
3609/// grammar, stated so the model writes conditions the plan validator accepts.
3610fn plan_instructions(min_nodes: u32) -> String {
3611    format!(
3612        " (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
3613\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
3614\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
3615with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
3616but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
3617'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
3618the cited evidence — and what the step does with it; and 'edges' from step to \
3619step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
3620'path != literal', 'path exists' or '!path', where path is dotted names and the \
3621literal is a JSON string, number, true, false or null — no other operators; state \
3622a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
3623earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
3624has branches or a loop; prefer a skill when it is a straight list. If a saved \
3625plan already covers this procedure, use ITS name so it is patched.",
3626        crate::llm::MAX_PLAN_NODES
3627    )
3628}
3629
3630/// What a `plan` proposal resolves to: the Skill grain (prose), the Workflow
3631/// grain (structure), the name both take from the target, the counts, and the
3632/// live pair of that name (to supersede) if there is one.
3633struct PlanFields {
3634    skill: serde_json::Map<String, Value>,
3635    /// `None` when the graph the model wrote would not run — an edge
3636    /// condition outside the runtime's frozen grammar, an edge naming no
3637    /// step. The procedure is still captured, as a skill.
3638    workflow: Option<serde_json::Map<String, Value>>,
3639    name: String,
3640    n_nodes: usize,
3641    n_edges: usize,
3642    existing_skill: Option<String>,
3643    existing_plan: Option<String>,
3644}
3645
3646/// The two grains a `plan` proposal would write.
3647///
3648/// Grounding is structural: every step's tool must be one the cited evidence
3649/// shows was called, so the model cannot plan around a tool it invented. The
3650/// workflow body is handed to the substrate's own plan validator before the
3651/// draft can be stamped applicable — a plan a reviewer could approve is one
3652/// the runtime would accept. The pair shares a `name`; the Workflow carries it
3653/// as a host field (the type is a container by design), and that is how the
3654/// live plan of a name is found to be patched rather than duplicated.
3655#[allow(clippy::too_many_arguments)]
3656fn derived_plan_fields<S: SubstrateRead>(
3657    sub: &S,
3658    target: &TargetRef,
3659    description: &str,
3660    when_to_use: &str,
3661    nodes: &[crate::llm::PlanNodeDraft],
3662    edges: &[crate::llm::PlanEdgeDraft],
3663    cited: &[String],
3664    ns_by_hash: &std::collections::BTreeMap<String, String>,
3665    plans: &crate::policy::PlanAuthoring,
3666) -> Option<PlanFields> {
3667    if target.scheme() != "entity" {
3668        return None;
3669    }
3670    let name = sanitize_skill_name(
3671        target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3672    )?;
3673    let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3674    let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3675    if description.is_empty() || when_to_use.is_empty() {
3676        return None;
3677    }
3678    if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
3679        return None;
3680    }
3681    // The tools the evidence shows were actually called.
3682    let known_tools: BTreeSet<String> = cited
3683        .iter()
3684        .filter_map(|h| sub.grain(h).ok().flatten())
3685        .filter_map(|g| g.tool_name().map(normalize_ident))
3686        .collect();
3687    let mut ids: Vec<String> = Vec::new();
3688    let mut steps: Vec<String> = Vec::new();
3689    let mut seen: BTreeSet<String> = BTreeSet::new();
3690    let mut grounded = 0usize;
3691    for n in nodes {
3692        let id = sanitize_skill_name(&n.id)?;
3693        if !seen.insert(id.clone()) {
3694            return None; // duplicate step id
3695        }
3696        let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
3697        if step.is_empty() {
3698            return None;
3699        }
3700        // A step whose tool the evidence never shows keeps its instruction and
3701        // loses the attribution — it is not recorded as calling anything. A
3702        // real procedure has steps that call nothing (deciding, grouping,
3703        // comparing), and a model writes them with a placeholder tool;
3704        // rejecting the whole draft for one of those threw away procedures
3705        // that were three-quarters grounded (PERSIST.md §11 #27). Nodes carry
3706        // no bindings here, so an unattributed step executes nothing and
3707        // claims nothing — but the prose must not tell a later session to
3708        // call a tool that does not exist.
3709        let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
3710        if !tool.is_empty() && known_tools.contains(&normalize_ident(&tool)) {
3711            grounded += 1;
3712            steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
3713        } else {
3714            steps.push(format!("{}. {id}: {step}", steps.len() + 1));
3715        }
3716        ids.push(id);
3717    }
3718    // Anchored in the trajectory: at least one step calls a tool the evidence
3719    // actually shows, or this is not a procedure the agent carried out.
3720    if grounded == 0 {
3721        return None;
3722    }
3723    // Edges are built leniently: one the runtime could not run costs the
3724    // plan, never the procedure. `runnable` goes false and the flow line is
3725    // still written into the skill's prose, where it is description rather
3726    // than a promise.
3727    let mut edge_vals: Vec<Value> = Vec::new();
3728    let mut flow_lines: Vec<String> = Vec::new();
3729    let mut runnable = true;
3730    for e in edges {
3731        let (Some(src), Some(dst)) = (sanitize_skill_name(&e.src), sanitize_skill_name(&e.dst)) else {
3732            runnable = false;
3733            continue;
3734        };
3735        if !seen.contains(&src) || !seen.contains(&dst) {
3736            // An edge naming no step — a model's "end" node, typically.
3737            flow_lines.push(format!("{} → {}", e.src.trim(), e.dst.trim()));
3738            runnable = false;
3739            continue;
3740        }
3741        let mut ev = serde_json::Map::new();
3742        ev.insert("src".into(), Value::from(src.clone()));
3743        ev.insert("dst".into(), Value::from(dst.clone()));
3744        let mut label = format!("{src} → {dst}");
3745        if let Some(c) = e
3746            .cond
3747            .as_deref()
3748            .map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
3749            .filter(|c| !c.is_empty())
3750        {
3751            label.push_str(&format!(" if {c}"));
3752            ev.insert("cond".into(), Value::from(c));
3753        }
3754        if let Some(m) = e.max_cycles {
3755            if m == 0 || m > 100 {
3756                runnable = false;
3757            } else {
3758                label.push_str(&format!(" (at most {m} times)"));
3759                ev.insert("max_cycles".into(), Value::from(m));
3760            }
3761        }
3762        flow_lines.push(label);
3763        edge_vals.push(Value::Object(ev));
3764    }
3765    if edge_vals.len() > 4 * ids.len() {
3766        runnable = false;
3767    }
3768    // Namespace: where the evidence lives, by majority — a lesson's rule.
3769    let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3770    for h in cited {
3771        if let Some(ns) = ns_by_hash.get(h) {
3772            if !ns.is_empty() {
3773                *ns_counts.entry(ns.as_str()).or_default() += 1;
3774            }
3775        }
3776    }
3777    let ns = ns_counts
3778        .iter()
3779        .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3780        .map(|(ns, _)| ns.to_string());
3781
3782    // The Workflow: what the runtime validates. Unbound steps are abstract
3783    // nodes — legal, and what a plan over a host's own tools is.
3784    let mut workflow = serde_json::Map::new();
3785    workflow.insert("nodes".into(), Value::from(ids.clone()));
3786    workflow.insert("edges".into(), Value::Array(edge_vals));
3787    workflow.insert("name".into(), Value::from(name.clone()));
3788    if let Some(ns) = &ns {
3789        workflow.insert("namespace".into(), Value::from(ns.clone()));
3790    }
3791    // The substrate owns the grammar: the runtime's own validator decides
3792    // whether this is a plan (unique and reachable steps, conditions that
3793    // parse, every cycle bounded). The engine carries no second opinion.
3794    let workflow = (runnable && sub.validate_plan(&Value::Object(workflow.clone())).is_ok())
3795        .then_some(workflow);
3796
3797    // The Skill: the same procedure as prose, with the graph's edges spelled
3798    // out under the steps so a reader sees the branches the plan encodes.
3799    let mut instructions = steps.join("\n");
3800    if !flow_lines.is_empty() {
3801        instructions.push_str("\n\nFlow:\n");
3802        instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
3803    }
3804    let mut skill = serde_json::Map::new();
3805    skill.insert("name".into(), Value::from(name.clone()));
3806    skill.insert("description".into(), Value::from(description));
3807    skill.insert("when_to_use".into(), Value::from(when_to_use));
3808    skill.insert("instructions".into(), Value::from(instructions));
3809    if let Some(ns) = &ns {
3810        skill.insert("namespace".into(), Value::from(ns.clone()));
3811    }
3812    let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
3813        sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
3814            .ok()?
3815            .into_iter()
3816            .find(|g| pick(g))
3817            .map(|g| g.hash)
3818    };
3819    let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
3820    let existing_plan = workflow
3821        .is_some()
3822        .then(|| live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str())))
3823        .flatten();
3824    let n_edges = workflow
3825        .as_ref()
3826        .and_then(|w| w.get("edges"))
3827        .and_then(Value::as_array)
3828        .map_or(0, |a| a.len());
3829    Some(PlanFields { skill, workflow, name, n_nodes: ids.len(), n_edges, existing_skill, existing_plan })
3830}
3831
3832/// The metric name under which the Verify gate records a premise that moved.
3833pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
3834
3835/// The Verify gate's second question. For every applied recommendation, the
3836/// grains it cited are looked up again: one that has been retracted, or
3837/// superseded by a grain holding a DIFFERENT value, is a premise that moved.
3838/// A value-identical supersession — what consolidation does — is not, and
3839/// neither is a supersession the recommendation's OWN apply performed: a
3840/// contradiction resolution cites the two conflicting facts and retires one
3841/// of them; that is the change it was approved to make, not its premise
3842/// moving out from under it.
3843///
3844/// Compared with the wrong rule (counting the apply's own work as drift), this
3845/// is what keeps the gate quiet on the analyzers that exist to supersede.
3846/// Records `drifted` in the outcome series once per distinct count (so a
3847/// pass does not re-record what the last pass already did) and returns an
3848/// input `outcome_review` turns into the revert proposal. The reviewer
3849/// decides; nothing here applies.
3850fn detect_premise_drift<S: OmsSubstrate>(
3851    sub: &S,
3852    p: &mut LoopPersisted,
3853    now_ms: i64,
3854) -> Result<Vec<OutcomeInput>> {
3855    let mut out = Vec::new();
3856    let applied: Vec<(String, String, Vec<String>)> = p
3857        .applied
3858        .iter()
3859        .filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
3860        .map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
3861        .collect();
3862    for (rec_hash, target_ref, own) in applied {
3863        let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
3864        if rec.evidence.is_empty() {
3865            continue;
3866        }
3867        let mut moved = 0u64;
3868        for e in &rec.evidence {
3869            match sub.grain(e)? {
3870                None => moved += 1, // retracted or gone
3871                Some(g) => {
3872                    let Some(newer) = &g.superseded_by else { continue };
3873                    if own.iter().any(|c| c == newer) {
3874                        continue; // the apply's own supersession
3875                    }
3876                    match sub.grain(newer)? {
3877                        // Superseded by something unreadable: a retraction
3878                        // (the reference substrate marks FORGET this way).
3879                        None => moved += 1,
3880                        Some(n) => {
3881                            if !same_value(&g, &n) {
3882                                moved += 1;
3883                            }
3884                        }
3885                    }
3886                }
3887            }
3888        }
3889        if moved == 0 {
3890            continue;
3891        }
3892        let already = p
3893            .outcomes
3894            .get(&rec_hash)
3895            .and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
3896            .is_some_and(|o| o.current == moved as f64);
3897        if !already {
3898            p.outcomes.entry(rec_hash.clone()).or_default().push(
3899                crate::recommendation::OutcomeResult {
3900                    rec_hash: rec_hash.clone(),
3901                    metric: PREMISE_DRIFT_METRIC.into(),
3902                    baseline: 0.0,
3903                    current: moved as f64,
3904                    verdict: "drifted".into(),
3905                    baseline_kind: "snapshot".into(),
3906                    baseline_run_id: None,
3907                    best_before: None,
3908                    tolerance: 0.0,
3909                    current_run_id: None,
3910                    cost: None,
3911                    horizon_ms: 0,
3912                    checkpoint: None,
3913                    measured_at_ms: now_ms,
3914                },
3915            );
3916        }
3917        out.push(OutcomeInput {
3918            rec_hash,
3919            target_ref,
3920            metric: PREMISE_DRIFT_METRIC.into(),
3921            baseline: 0.0,
3922            current: moved as f64,
3923            unit: "superseded premises".into(),
3924            higher_is_better: false,
3925            baseline_kind: "snapshot".into(),
3926            baseline_run_id: None,
3927            best_before: None,
3928            tolerance: 0.0,
3929            current_run_id: None,
3930            cost: None,
3931        });
3932    }
3933    Ok(out)
3934}
3935
3936/// Does the superseding grain say the same thing as the one it replaced? A
3937/// fact compares its object; anything else compares its text body. Two grains
3938/// that cannot be compared are treated as different — the fail-closed
3939/// reading, since a premise we cannot confirm still holds is one that moved.
3940fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
3941    if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
3942        return normalize_ident(a) == normalize_ident(b);
3943    }
3944    for key in ["content", "tool_content", "body", "text", "object"] {
3945        if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
3946            return normalize_ident(a) == normalize_ident(b);
3947        }
3948    }
3949    false
3950}
3951
3952/// A skill name is an identifier: `[A-Za-z0-9_-]`, bounded, case preserved.
3953fn sanitize_skill_name(s: &str) -> Option<String> {
3954    let t = s.trim();
3955    if t.is_empty()
3956        || t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
3957        || !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
3958    {
3959        return None;
3960    }
3961    Some(t.to_string())
3962}
3963
3964/// What a `skill` proposal resolves to: the grain's fields, the name it took
3965/// from the target, how many steps survived sanitizing, and the live skill of
3966/// that name in the same namespace (to supersede) if there is one.
3967struct SkillFields {
3968    fields: serde_json::Map<String, Value>,
3969    name: String,
3970    n_steps: usize,
3971    existing: Option<String>,
3972}
3973
3974/// The fields of the Skill grain a `skill` proposal would write.
3975///
3976/// The namespace is the one most of the cited evidence lives in — the same
3977/// rule a lesson follows — never one the model names.
3978#[allow(clippy::too_many_arguments)]
3979fn derived_skill_fields<S: SubstrateRead>(
3980    sub: &S,
3981    target: &TargetRef,
3982    description: &str,
3983    when_to_use: &str,
3984    steps: &[String],
3985    cited: &[String],
3986    ns_by_hash: &std::collections::BTreeMap<String, String>,
3987    skills: &crate::policy::SkillAuthoring,
3988) -> Option<SkillFields> {
3989    if target.scheme() != "entity" {
3990        return None;
3991    }
3992    let name = sanitize_skill_name(
3993        target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3994    )?;
3995    let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3996    let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3997    let steps: Vec<String> = steps
3998        .iter()
3999        .map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
4000        .filter(|st| !st.is_empty())
4001        .take(crate::llm::MAX_SKILL_STEPS)
4002        .collect();
4003    if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
4004        return None;
4005    }
4006    // Namespace: where the evidence lives, by majority (ties → lexically
4007    // first), exactly as a lesson's.
4008    let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
4009    for h in cited {
4010        if let Some(ns) = ns_by_hash.get(h) {
4011            if !ns.is_empty() {
4012                *ns_counts.entry(ns.as_str()).or_default() += 1;
4013            }
4014        }
4015    }
4016    let ns = ns_counts
4017        .iter()
4018        .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
4019        .map(|(ns, _)| ns.to_string());
4020    let instructions = steps
4021        .iter()
4022        .enumerate()
4023        .map(|(i, st)| format!("{}. {st}", i + 1))
4024        .collect::<Vec<_>>()
4025        .join("\n");
4026    let mut fields = serde_json::Map::new();
4027    fields.insert("name".into(), Value::from(name.clone()));
4028    fields.insert("description".into(), Value::from(description));
4029    fields.insert("when_to_use".into(), Value::from(when_to_use));
4030    fields.insert("instructions".into(), Value::from(instructions));
4031    if let Some(ns) = &ns {
4032        fields.insert("namespace".into(), Value::from(ns.clone()));
4033    }
4034    // Patch, don't duplicate: the live skill of this name in this namespace.
4035    let existing = sub
4036        .grains_of_type(
4037            crate::model::grain_type::SKILL,
4038            ns.as_deref(),
4039            ReadOpts { live_only: true, since_ms: None },
4040        )
4041        .ok()?
4042        .into_iter()
4043        .find(|g| g.skill_name() == Some(name.as_str()))
4044        .map(|g| g.hash);
4045    Some(SkillFields { fields, name, n_steps: steps.len(), existing })
4046}
4047
4048fn derived_fact_fields(
4049    target: &TargetRef,
4050    relation: &str,
4051    object: &str,
4052    cited: &[String],
4053    ns_by_hash: &std::collections::BTreeMap<String, String>,
4054) -> Option<serde_json::Map<String, Value>> {
4055    if target.scheme() != "entity" {
4056        return None;
4057    }
4058    let subject = target
4059        .opaque()
4060        .rsplit_once('/')
4061        .map(|(_, s)| s)
4062        .unwrap_or(target.opaque());
4063    if subject.is_empty() {
4064        return None;
4065    }
4066    let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
4067    for h in cited {
4068        if let Some(ns) = ns_by_hash.get(h) {
4069            if !ns.is_empty() {
4070                *ns_counts.entry(ns.as_str()).or_default() += 1;
4071            }
4072        }
4073    }
4074    let lesson_ns = ns_counts
4075        .iter()
4076        .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
4077        .map(|(ns, _)| ns.to_string());
4078    let mut fields = serde_json::Map::new();
4079    fields.insert("subject".into(), Value::from(subject));
4080    fields.insert("relation".into(), Value::from(relation));
4081    fields.insert("object".into(), Value::from(object));
4082    // `confidence` is stamped by the caller from the VERIFIER's calibrated
4083    // score, not the proposer's self-report — so it is deliberately absent
4084    // here, where only the proposer has spoken.
4085    if let Some(ns) = lesson_ns {
4086        fields.insert("namespace".into(), Value::from(ns));
4087    }
4088    Some(fields)
4089}
4090
4091/// The latest Verify-gate verdict for every grain an applied recommendation
4092/// created — keyed by the CREATED hash, so an analyzer looking at a live
4093/// lesson can say how it measured (`held`, `regressed`, `drifted`,
4094/// `held_costlier`) without reaching the engine's state itself.
4095fn latest_verdicts(p: &LoopPersisted) -> BTreeMap<String, String> {
4096    let mut out = BTreeMap::new();
4097    for (rec_hash, applied) in &p.applied {
4098        let latest = p
4099            .outcomes
4100            .get(rec_hash)
4101            .and_then(|v| v.iter().max_by_key(|o| o.measured_at_ms))
4102            .map(|o| o.verdict.clone());
4103        for h in &applied.created_hashes {
4104            out.insert(h.clone(), latest.clone().unwrap_or_else(|| "unmeasured".into()));
4105        }
4106    }
4107    out
4108}
4109
4110/// Cosine similarity at or above which two lesson embeddings are one
4111/// instruction (the T1 leg).
4112pub const NEAR_DUPLICATE_COSINE: f64 = 0.90;
4113/// Token-set Jaccard at or above which two lesson texts are one instruction
4114/// (the T0 floor — weak, and honest about it).
4115pub const NEAR_DUPLICATE_JACCARD: f64 = 0.60;
4116/// How many near-duplicates a recommendation names, best first.
4117const NEAR_DUPLICATE_CAP: usize = 8;
4118
4119/// The live lessons on `subject` (in `namespace`, when known) that `text`
4120/// restates. Cosine over the substrate's embedder when it embeds both sides;
4121/// otherwise normalized token-set Jaccard. Best first, capped. A read that
4122/// fails yields nothing — a near-duplicate check must never block a draft.
4123pub(crate) fn near_duplicates_of<S: SubstrateRead + ?Sized>(
4124    sub: &S,
4125    subject: &str,
4126    namespace: Option<&str>,
4127    text: &str,
4128) -> Vec<crate::recommendation::NearDuplicate> {
4129    use crate::analyzers::duplicate_sweep::{jaccard, tokenize};
4130    if subject.is_empty() || text.trim().is_empty() {
4131        return Vec::new();
4132    }
4133    let Ok(facts) = sub.grains_of_type(
4134        crate::model::grain_type::FACT,
4135        namespace,
4136        ReadOpts { live_only: true, since_ms: None },
4137    ) else {
4138        return Vec::new();
4139    };
4140    let mine = sub.embed(text).ok().flatten();
4141    let my_tokens = tokenize(text);
4142    let mut out: Vec<crate::recommendation::NearDuplicate> = facts
4143        .iter()
4144        .filter(|f| f.fact_relation() == Some("lesson"))
4145        .filter(|f| f.fact_subject().is_some_and(|s| normalize_ident(s) == normalize_ident(subject)))
4146        .filter_map(|f| {
4147            let other = f.fact_object()?;
4148            let (score, method, floor) = match (&mine, sub.embed(other).ok().flatten()) {
4149                (Some(a), Some(b)) => (cosine(a, &b), "cosine", NEAR_DUPLICATE_COSINE),
4150                _ => (jaccard(&my_tokens, &tokenize(other)), "jaccard", NEAR_DUPLICATE_JACCARD),
4151            };
4152            (score >= floor).then(|| crate::recommendation::NearDuplicate {
4153                hash: f.hash.clone(),
4154                score: (score * 1000.0).round() / 1000.0,
4155                method: method.into(),
4156            })
4157        })
4158        .collect();
4159    out.sort_by(|a, b| b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal).then(a.hash.cmp(&b.hash)));
4160    out.truncate(NEAR_DUPLICATE_CAP);
4161    out
4162}
4163
4164fn cosine(a: &[f32], b: &[f32]) -> f64 {
4165    if a.len() != b.len() || a.is_empty() {
4166        return 0.0;
4167    }
4168    let (mut dot, mut na, mut nb) = (0f64, 0f64, 0f64);
4169    for (x, y) in a.iter().zip(b) {
4170        dot += *x as f64 * *y as f64;
4171        na += *x as f64 * *x as f64;
4172        nb += *y as f64 * *y as f64;
4173    }
4174    if na == 0.0 || nb == 0.0 {
4175        0.0
4176    } else {
4177        dot / (na.sqrt() * nb.sqrt())
4178    }
4179}
4180
4181/// The `consolidation` paragraph appended to the DISCOVER instructions. It
4182/// is answerable only to a `lesson_pile` finding, which lists the hashes the
4183/// proposal must name — the model cannot pick a pile of its own.
4184const CONSOLIDATION_INSTRUCTIONS: &str = " (8) {\"kind\":\"consolidation\",\"lesson\":\"...\",\
4185\"supersedes\":[\"<hash>\",...]} with the same entity target — ONLY in answer to a \
4186'Lesson pile' finding, which lists the live lessons on one entity that exceed \
4187its budget. Write ONE short imperative rule (max 240 chars) that says what \
4188those lessons say together, dropping nothing a lesson that measured 'held' \
4189required and keeping nothing only a lesson that measured 'regressed' or \
4190'drifted' added; 'supersedes' MUST be exactly the hashes that finding lists \
4191(cite them as evidence too). Applying it replaces every listed lesson with \
4192the one line; the reviewer can restore them all.";
4193
4194/// The action kinds that apply ONLY through the evalset-run gating edge
4195/// (§7.4 for tool code; the tuning seam's adapter promotion inherits the
4196/// same rule). One predicate so the gate, the rollbackable stamp, and the
4197/// promotion write can never disagree on membership.
4198fn requires_gating(kind: ActionKind) -> bool {
4199    matches!(
4200        kind,
4201        ActionKind::CodeRevision | ActionKind::AdapterRevision
4202    )
4203}
4204
4205fn stamp(
4206    m: &AnalyzerManifest,
4207    params: &crate::manifest::Params,
4208    d: crate::recommendation::RecDraft,
4209    now_ms: i64,
4210) -> Result<Recommendation> {
4211    let target = TargetRef::parse(&d.target_ref)?;
4212    // Rule E1 at the door (§7.4): an unpinned code revision, a pinned
4213    // non-code target, or a code action on a non-tool target never becomes
4214    // a recommendation at all.
4215    crate::recommendation::validate_code_rules(
4216        d.action_kind,
4217        target.target_class(),
4218        d.evalset_hash.as_deref(),
4219    )?;
4220    // A revert's identity is the recommendation it retracts, not just its
4221    // target: two regressed lessons on one entity are two reverts.
4222    let revert_of = match (&d.action_kind, &d.proposal) {
4223        (ActionKind::Revert, Proposal::Data { data }) => {
4224            data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
4225        }
4226        _ => None,
4227    };
4228    let dedup = match revert_of.as_deref() {
4229        Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
4230        None => dedup_key(m.family(), &d.target_ref, d.action_kind),
4231    };
4232    let destructive = match &d.proposal {
4233        Proposal::Cal { cal } => cal::contains_destructive(cal),
4234        _ => false,
4235    };
4236    let rollbackable = match &d.proposal {
4237        Proposal::Cal { .. } => !destructive,
4238        Proposal::Edit { .. } => false,
4239        // A code or adapter revision applies by WRITING the promotion
4240        // grain; retracting it is the exact inverse — rollbackable by
4241        // construction.
4242        Proposal::Data { .. } => requires_gating(d.action_kind),
4243    };
4244    let mut evidence = d.evidence;
4245    evidence.truncate(MAX_EVIDENCE);
4246    // Provenance follows the analyzer's trust class: a subprocess
4247    // (`--analyzer-cmd`) finding is stamped `Command` — structurally
4248    // auto-apply-ineligible and badged [external] on the recall surface —
4249    // not `Builtin`.
4250    let origin = match m.trust_class {
4251        crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
4252        _ => Origin::Builtin,
4253    };
4254    Ok(Recommendation {
4255        hash: String::new(),
4256        analyzer: m.id.clone(),
4257        params_snapshot: params.snapshot(),
4258        origin,
4259        target_ref: target.as_string(),
4260        action_kind: d.action_kind,
4261        dedup_key: dedup,
4262        summary: d.summary,
4263        severity: d.severity,
4264        proposal: d.proposal,
4265        destructive,
4266        rollbackable,
4267        evidence,
4268        evidence_query: d.evidence_query,
4269        metric: d.metric,
4270        confidence: d.confidence,
4271        importance: d.importance,
4272        created_at_ms: now_ms,
4273        guidance: None,
4274        evalset_hash: d.evalset_hash,
4275        near_duplicate_of: Vec::new(),
4276        replay: None,
4277        status: RecStatus::Pending,
4278    })
4279}
4280
4281fn validate_because(because: &str) -> Result<String> {
4282    let trimmed = because.trim();
4283    if trimmed.is_empty() {
4284        return Err(Error::InvalidProposal(
4285            "a BECAUSE reason is required".into(),
4286        ));
4287    }
4288    if trimmed.chars().count() > MAX_BECAUSE {
4289        return Err(Error::InvalidProposal(format!(
4290            "BECAUSE exceeds {MAX_BECAUSE} chars"
4291        )));
4292    }
4293    Ok(trimmed.to_string())
4294}
4295
4296fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
4297    for req in &m.requires {
4298        match req {
4299            Capability::Forks if !caps.forks => return Some("forks"),
4300            Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
4301            Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
4302            _ => {}
4303        }
4304    }
4305    None
4306}
4307
4308fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
4309    p.config.get(analyzer_id).and_then(|c| c.severity_floor)
4310}
4311
4312fn gate(
4313    opts: &RunOptions,
4314    p: &LoopPersisted,
4315    new_grains: u64,
4316    new_errors: u64,
4317    now_ms: i64,
4318) -> Option<SkipReason> {
4319    let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
4320    if !any {
4321        return None;
4322    }
4323    let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
4324    let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
4325    let stale_ok = opts
4326        .if_stale_ms
4327        .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
4328    if min_new_ok || min_err_ok || stale_ok {
4329        return None;
4330    }
4331    // Choose the honest reason: staleness-only gate → not_stale, else min_new.
4332    if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
4333        Some(SkipReason::NotStale)
4334    } else {
4335        Some(SkipReason::MinNewNotMet)
4336    }
4337}
4338
4339/// What landed since a watermark, in every unit a gate can count.
4340#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
4341pub(crate) struct NewSince {
4342    /// Grains of the four evidence-bearing types.
4343    pub grains: u64,
4344    /// Tool grains recording a failure (`--min-new-errors`).
4345    pub error_events: u64,
4346    /// Event grains — turns, in a chat deployment (`cadence.every_events`).
4347    pub events: u64,
4348    /// Distinct `session_id`s among those Events (`cadence.every_sessions`).
4349    pub sessions: u64,
4350}
4351
4352fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
4353    let opts = ReadOpts {
4354        live_only: false,
4355        since_ms: watermark.map(|w| w + 1),
4356    };
4357    let mut n = NewSince::default();
4358    let mut sessions: BTreeSet<&str> = BTreeSet::new();
4359    let mut events_held: Vec<GrainRecord> = Vec::new();
4360    for t in [
4361        crate::model::grain_type::FACT,
4362        crate::model::grain_type::EVENT,
4363        crate::model::grain_type::TOOL,
4364        crate::model::grain_type::OBSERVATION,
4365    ] {
4366        let g = sub.grains_of_type(t, None, opts)?;
4367        n.grains += g.len() as u64;
4368        // The error gate (--min-new-errors) watches captured tool failures.
4369        if t == crate::model::grain_type::TOOL {
4370            n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
4371        }
4372        if t == crate::model::grain_type::EVENT {
4373            n.events = g.len() as u64;
4374            events_held = g;
4375        }
4376    }
4377    for e in &events_held {
4378        if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
4379            sessions.insert(sid);
4380        }
4381    }
4382    n.sessions = sessions.len() as u64;
4383    Ok(n)
4384}
4385
4386/// The policy cadence, evaluated: `None` when a pass is due.
4387///
4388/// OR over the thresholds the host set — the pass runs when any one is met.
4389/// The same shape as the per-call gate above, deliberately: the flags and
4390/// the policy block are one mechanism spelled in two places, and the flags
4391/// win when both are present.
4392fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
4393    if !c.is_set() {
4394        return None;
4395    }
4396    let time_ok = c
4397        .every_ms
4398        .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
4399    let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
4400    let events_ok = c.every_events.is_some_and(|m| new.events >= m);
4401    let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
4402    if time_ok || grains_ok || events_ok || sessions_ok {
4403        None
4404    } else {
4405        Some(SkipReason::CadenceNotDue)
4406    }
4407}
4408
4409fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
4410    let grains = sub.grains_of_type(
4411        crate::model::grain_type::RECOMMENDATION,
4412        Some(LOOP_NS),
4413        ReadOpts {
4414            live_only: false,
4415            since_ms: None,
4416        },
4417    )?;
4418    let mut set = BTreeSet::new();
4419    for g in grains {
4420        let status = p
4421            .status_index
4422            .get(&g.hash)
4423            .copied()
4424            .unwrap_or(RecStatus::Pending);
4425        // Pending/approved (still open) and applied (already handled)
4426        // recommendations suppress re-proposal of the same finding. Rejected
4427        // is handled by cooldowns, and so is a rollback the Verify gate
4428        // caused (`strike_cooldown` at the revert apply); an operator's own
4429        // rollback and expiry may legitimately re-propose (the situation
4430        // returned).
4431        if matches!(
4432            status,
4433            RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
4434        ) {
4435            if let Some(key) = g.str_field("dedup_key") {
4436                set.insert(key.to_string());
4437            }
4438        }
4439    }
4440    Ok(set)
4441}
4442
4443/// Put a finding's `dedup_key` on an exponential cooldown: 7d, 14d, 28d, …
4444/// capped at 90d, so a finding a reviewer keeps rejecting stops re-surfacing
4445/// on a fixed 7d cadence (it was a flat 7d despite the "doubling" comment).
4446/// Two events earn a strike: a reviewer's rejection, and a revert the Verify
4447/// gate proposed on a measured regression — both are a verdict that the
4448/// finding, as it stands, should not come back on the next pass.
4449pub(crate) fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
4450    const BASE_MS: i64 = 7 * 86_400_000;
4451    const CAP_MS: i64 = 90 * 86_400_000;
4452    let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
4453    let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
4454    *strikes = strikes.saturating_add(1);
4455    p.cooldowns.insert(dedup_key, now_ms + interval);
4456}
4457
4458fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
4459    let g = sub
4460        .grain(rec_hash)?
4461        .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
4462    Recommendation::from_fields(rec_hash, &g.fields)
4463}
4464
4465/// Is this line a definition rewrite — a statement that changes a saved
4466/// `qry:`/`tpl:` registry row rather than writing a grain?
4467///
4468/// A keyword test, not a parse: the engine deliberately contains a CAL
4469/// *writer*, never a parser (parsing is the substrate's job). Both spellings
4470/// are matched case-insensitively, and `DROP` is intentionally absent — the
4471/// loop may propose defining a query, never removing one.
4472pub(crate) fn is_definition_statement(line: &str) -> bool {
4473    let up = line.trim_start().to_ascii_uppercase();
4474    up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
4475}
4476
4477/// The refusal an advisory `Edit` earns. The engine has no executable edit
4478/// primitive; the change belongs in the host.
4479const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
4480     Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
4481     approve it to acknowledge it and let it expire.";
4482
4483/// The refusal an advisory `Data` finding earns — every `Data` shape except
4484/// `outcome_review`'s revert, which carries `revert_of`.
4485/// Shared by [`Engine::preflight_apply`] and the apply gate so a fused
4486/// approve-and-apply caller is refused BEFORE the approval lands, not after.
4487const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
4488     (evalset hash + run id + stats) — use apply_gated";
4489
4490const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
4491     guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
4492     acknowledge it and let it expire.";
4493
4494/// Whether [`Engine::apply`] can execute this proposal at all.
4495///
4496/// One source of truth, shared by [`Engine::preflight_apply`] and
4497/// [`Engine::apply`] so the two can never disagree.
4498///
4499/// Deliberately **not** consulted on approve. An advisory finding is a `Flag`:
4500/// approving it means "yes, this is real", which is the whole workflow for the
4501/// LLM path and the telemetry analyzers. What must not happen is a caller being
4502/// walked into an approval and *then* refused — which is exactly what the fused
4503/// approve-and-apply path in the bindings did, leaving the recommendation in
4504/// `approved`, whose only exits are `applied` and `expired`. Preflight asks
4505/// first, so that path now refuses before it commits anything.
4506pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
4507    match proposal {
4508        Proposal::Cal { .. } => Ok(()),
4509        Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
4510        // Two executable Data shapes: `outcome_review`'s `revert_of` (names
4511        // an earlier applied recommendation to roll back), and a gated
4512        // revision (code or adapter), which executes by writing its
4513        // promotion grain — the gating-run requirement itself is checked at
4514        // apply, not here.
4515        Proposal::Data { data } => {
4516            if requires_gating(action_kind)
4517                || data.get("revert_of").and_then(Value::as_str).is_some()
4518            {
4519                Ok(())
4520            } else {
4521                Err(Error::InvalidProposal(ADVISORY_DATA.into()))
4522            }
4523        }
4524    }
4525}
4526
4527#[cfg(test)]
4528mod definition_body_tests {
4529    use super::safe_definition_body;
4530
4531    #[test]
4532    fn ordinary_bodies_pass() {
4533        assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
4534        assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
4535    }
4536
4537    #[test]
4538    fn a_body_cannot_close_its_own_block_or_carry_destruction() {
4539        // The injection shape: close the DEFINE block, append a statement.
4540        // Newlines are already collapsed by `sanitize_line`, so the payload
4541        // arrives as ONE line — which is exactly what a line-leading-keyword
4542        // destructive scan cannot see.
4543        assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
4544        assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
4545        // …and the keyword alone is refused even without the braces.
4546        assert!(!safe_definition_body("RECALL facts FORGET abc"));
4547        assert!(!safe_definition_body("recall facts purge older than 1d"));
4548        assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
4549        // A nested DEFINE would redefine something the target does not name.
4550        assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
4551        // Substring matches are not keywords — this must still pass.
4552        assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
4553    }
4554}
4555
4556#[cfg(test)]
4557mod plan_edit_tests {
4558    use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
4559    use serde_json::json;
4560
4561    fn plan() -> serde_json::Value {
4562        json!({
4563            "nodes": ["fetch", "review", "post"],
4564            "edges": [
4565                {"src": "fetch", "dst": "review"},
4566                {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
4567            ],
4568            "bindings": {"fetch": "sha256:tool1"},
4569            "retries": {"fetch": 1}
4570        })
4571    }
4572
4573    #[test]
4574    fn the_allowlist_admits_thresholds_and_refuses_topology() {
4575        assert!(plan_edit_allowed("edges.1.cond"));
4576        assert!(plan_edit_allowed("edges.1.max_cycles"));
4577        assert!(plan_edit_allowed("retries.fetch"));
4578        // Topology is not expressible — the structural half of the guarantee
4579        // that a plan revision stays reviewable as scalar deltas.
4580        for path in [
4581            "nodes",
4582            "nodes.0",
4583            "edges.0.src",
4584            "edges.0.dst",
4585            "edges",
4586            "bindings.fetch",
4587            "edges.x.cond",
4588            "",
4589        ] {
4590            assert!(!plan_edit_allowed(path), "{path} must not be editable");
4591        }
4592    }
4593
4594    #[test]
4595    fn values_are_type_checked_against_the_field() {
4596        // A string in max_cycles would be DROPPED by the grain deserializer,
4597        // so an "applied" tightening would silently mean unlimited.
4598        assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
4599        assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
4600        assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
4601        assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
4602        assert!(plan_value_ok("retries.fetch", &json!(3)));
4603        assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
4604        assert!(!plan_value_ok("edges.1.cond", &json!("  ")));
4605        assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
4606        assert!(!plan_value_ok("edges.0.src", &json!("other")));
4607    }
4608
4609    #[test]
4610    fn get_reads_through_arrays_and_objects_and_absence_is_null() {
4611        let p = plan();
4612        assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
4613        assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
4614        // An edit that ADDS a retry declares `from: null` — so absence has to
4615        // read as Null rather than as an error.
4616        assert_eq!(plan_get(&p, "retries.review"), json!(null));
4617        assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
4618        assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
4619    }
4620
4621    #[test]
4622    fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
4623        let mut p = plan();
4624        assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
4625        assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
4626        assert!(plan_set(&mut p, "retries.review", json!(2)));
4627        assert_eq!(plan_get(&p, "retries.review"), json!(2));
4628        assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
4629        assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
4630        // A stored plan with no retries omits the map entirely (omit-defaults
4631        // serialization); the first retry edit on it must create the map, or
4632        // the revision can never resolve. Found by the golden E2E on the demo
4633        // plan, which the reference fixture — it always had a map — hid.
4634        let mut q = plan();
4635        q.as_object_mut().unwrap().remove("retries");
4636        assert!(plan_set(&mut q, "retries.greet", json!(1)));
4637        assert_eq!(plan_get(&q, "retries.greet"), json!(1));
4638        // Only `retries` is created; any other missing parent still refuses.
4639        let mut r = plan();
4640        r.as_object_mut().unwrap().remove("edges");
4641        assert!(!plan_set(&mut r, "edges.0.cond", json!("x")));
4642    }
4643}
4644
4645#[cfg(test)]
4646mod definition_proposal_tests {
4647    use super::is_definition_statement;
4648
4649    #[test]
4650    fn definition_statements_are_recognized_in_both_spellings() {
4651        assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
4652        assert!(is_definition_statement("  define template foo AS { x }"));
4653        assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
4654        // Ordinary proposals are untouched.
4655        assert!(!is_definition_statement("ADD fact {}"));
4656        assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
4657        assert!(!is_definition_statement("FORGET abc"));
4658        // `DROP` is never proposable, so it is deliberately NOT a definition
4659        // statement here — a proposal containing one still fails validation
4660        // rather than being handed an inverse.
4661        assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
4662    }
4663}