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