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