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