Skip to main content

areev_loop/
engine.rs

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