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, or an instruction from a \
2331person is ALSO penalized 1. Abstain only when the evidence shows none of \
2332those. Prefer the one proposal that addresses the most frequent or most costly \
2333failure over several speculative ones, and report your confidence honestly — \
2334an independent verifier, not you, decides what survives."
2335);
2336
2337/// The fixed GROUND instruction (§5.2): verify the finding's factual PREMISES
2338/// are real (anti-fabrication), while allowing an inference. A self-improvement
2339/// finding may reason BEYOND the evidence (e.g. 'HQ=SF and country=Germany are
2340/// inconsistent'); grounding checks the premises (HQ=SF, country=Germany) are in
2341/// the evidence — the soundness of the inference is VERIFY's job, not this one.
2342const GROUND_INSTRUCTIONS: &str = "You are a grounding checker guarding against \
2343fabrication. A finding may draw an INFERENCE from facts — your job is to confirm \
2344the facts it relies on are actually present in the cited evidence, NOT that its \
2345conclusion is stated verbatim. Decompose the finding into the factual claims it \
2346depends on. Mark supported=true when those facts are present in the evidence \
2347(even if the finding reasons beyond them). Mark supported=false ONLY if it relies \
2348on a fact that is NOT in the evidence (a fabrication) or cites evidence about a \
2349different subject. Return JSON: {\"results\":[{\"id\":0,\"supported\":true,\"reason\":\"...\"}]}.";
2350
2351/// The fixed VERIFY instruction (§5.3): adversarial, abstention-biased.
2352const VERIFY_INSTRUCTIONS: &str = "You are an adversarial reviewer stress-testing \
2353each finding for SOUNDNESS — reject unsound findings, but only for real reasons, \
2354never invented ones. For each finding, ask: (1) Reality — is there a genuine, \
2355SPECIFIC problem, or is it vague/speculative? Reject hedged 'potential' or \
2356'possible' findings with no concrete defect, and reject any claimed \
2357inconsistency or contradiction that is not backed by at least two actually \
2358conflicting facts in the cited evidence. (2) Context — does the finding \
2359correctly read its cited evidence, or misinterpret what the grains say? Keep a \
2360finding when it names a genuine, specific problem grounded in its evidence and \
2361materially useful to a human reviewer; otherwise reject it, and default to \
2362keep=false when uncertain. Do NOT reject a finding for being 'already known', \
2363redundant, or a 'common type of error' — duplication is handled elsewhere, and a \
2364grounded cross-fact inconsistency with two conflicting facts is exactly what to \
2365KEEP. Give a calibrated confidence 0.0-1.0. Return JSON: \
2366{\"results\":[{\"id\":0,\"keep\":true,\"confidence\":0.0,\"reason\":\"...\"}]}.";
2367
2368/// The fixed ENRICH instruction.
2369const ENRICH_INSTRUCTIONS: &str = "For each finding, optionally add a one-sentence \
2370guidance note to help a human reviewer decide. Do not restate the finding. Return \
2371JSON: {\"notes\":[{\"target\":\"<target_ref>\",\"guidance\":\"...\"}]}.";
2372
2373/// Add a grain to the evidence bundle — deduplicated, bounded at 64, text
2374/// capped. Shared by the deterministic-citation and recent-grain seeding.
2375/// `ns_by_hash` records each bundled grain's namespace so an authored lesson
2376/// can later land in the namespace its evidence lives in (never one the
2377/// model names).
2378fn push_evidence(
2379    evidence: &mut Vec<crate::llm::EvidenceItem>,
2380    bundle: &mut BTreeSet<String>,
2381    ns_by_hash: &mut std::collections::BTreeMap<String, String>,
2382    g: &GrainRecord,
2383    attribution: crate::policy::EvidenceAttribution,
2384) {
2385    if evidence.len() < EVIDENCE_CAP && bundle.insert(g.hash.clone()) {
2386        ns_by_hash.insert(g.hash.clone(), g.namespace.clone());
2387        evidence.push(crate::llm::EvidenceItem {
2388            id: format!("e{}", evidence.len() + 1),
2389            hash: g.hash.clone(),
2390            grain_type: g.grain_type.clone(),
2391            text: crate::llm::cap(&grain_brief_with(g, attribution), 400),
2392        });
2393    }
2394}
2395
2396/// Resolve one citation to a bundled grain's hash: the full hash, the
2397/// bundle-local `id` the evidence item carried, or an unambiguous hash prefix
2398/// of at least 12 hex chars. Anything else is a fabrication and resolves to
2399/// nothing. Small models copy 64-hex hashes badly — measured live, the
2400/// cite-check was where most of a cheap model's drafts died ("proposed 3 →
2401/// cited 1") — and a citation that plainly names one bundled grain is not
2402/// the thing that check exists to catch.
2403pub(crate) fn resolve_citation(
2404    cite: &str,
2405    bundle: &BTreeSet<String>,
2406    id_to_hash: &std::collections::BTreeMap<&str, &str>,
2407) -> Option<String> {
2408    let cite = cite.trim();
2409    if bundle.contains(cite) {
2410        return Some(cite.to_string());
2411    }
2412    if let Some(h) = id_to_hash.get(cite) {
2413        return Some((*h).to_string());
2414    }
2415    const MIN_PREFIX: usize = 12;
2416    if cite.len() >= MIN_PREFIX && cite.chars().all(|c| c.is_ascii_hexdigit()) {
2417        let lower = cite.to_ascii_lowercase();
2418        let mut it = bundle.iter().filter(|h| h.starts_with(&lower));
2419        if let (Some(h), None) = (it.next(), it.next()) {
2420            return Some(h.clone());
2421        }
2422    }
2423    None
2424}
2425
2426/// A short human-readable projection of a grain for the evidence bundle,
2427/// under the host's attribution policy.
2428fn grain_brief_with(g: &GrainRecord, attribution: crate::policy::EvidenceAttribution) -> String {
2429    if let (Some(s), Some(r), Some(o)) = (g.fact_subject(), g.fact_relation(), g.fact_object()) {
2430        return format!("{s} {r} {o}");
2431    }
2432    // Tool grains — the evidence most lesson drafts cite. Rendering them
2433    // empty starved GROUND of the very facts it exists to check: a correct
2434    // lesson would be refused as unverifiable (found live — the gate
2435    // rightly rejected a claim over evidence it could not see).
2436    if let Some(t) = g.tool_name() {
2437        let status = if g.is_error() { "error" } else { "ok" };
2438        let out = g.tool_content().unwrap_or("");
2439        // The call's input, when recorded: without it a successful trajectory
2440        // reads as a list of tool names and outputs, and a procedure — WHICH
2441        // ticket was fetched, WHAT tag was set — cannot be reconstructed from
2442        // it. A skill proposal needs the arguments; a lesson usually does not,
2443        // and the cap on the brief bounds the cost either way.
2444        let input = match g.fields.get("input") {
2445            Some(Value::String(v)) if !v.is_empty() => format!(" input={v}"),
2446            Some(v @ Value::Object(_)) | Some(v @ Value::Array(_)) => format!(" input={v}"),
2447            _ => String::new(),
2448        };
2449        return format!("tool {t}{input} {status}: {out}");
2450    }
2451    // `object` is last but it is not optional: an Observation stores its text
2452    // there (subject + object, no relation), so it misses the fact-triple
2453    // branch above and used to fall through this list to an empty string —
2454    // every human note in a memory reached the model as a blank line. That is
2455    // the single highest-value evidence a memory holds, and it was the one
2456    // shape that rendered to nothing. Callers had started duplicating the text
2457    // into `body` to work around it; nothing should have to.
2458    for key in ["content", "body", "text", "summary", "object"] {
2459        if let Some(v) = g.fields.get(key).and_then(|v| v.as_str()) {
2460            if v.is_empty() {
2461                continue;
2462            }
2463            // Who said it, when the grain records it. An Observation reaches
2464            // the model as a bare sentence otherwise, and a bare sentence is
2465            // ambiguous about direction in exactly the way that matters: a
2466            // person's correction ("Vendor Name is ACME") reads identically
2467            // to the agent having been told something it asked for. Measured
2468            // live on the receipts corpus, a model given 31 unattributed
2469            // corrections concluded the agent was repeatedly *requesting*
2470            // data it already had, and proposed rules to stop it asking.
2471            // The observer is already on the grain; only the projection
2472            // dropped it.
2473            if attribution == crate::policy::EvidenceAttribution::Anonymous {
2474                return v.to_string();
2475            }
2476            if let Some(who) = g.fields.get("observer_id").and_then(|v| v.as_str()) {
2477                if !who.is_empty() {
2478                    let kind = g
2479                        .fields
2480                        .get("observer_type")
2481                        .and_then(|v| v.as_str())
2482                        .unwrap_or("");
2483                    let about = g.fields.get("subject").and_then(|v| v.as_str()).unwrap_or("");
2484                    let mut prefix = if kind == "human" {
2485                        format!("{who} (a person) said")
2486                    } else {
2487                        format!("{who} observed")
2488                    };
2489                    if !about.is_empty() {
2490                        prefix.push_str(&format!(" of {about}"));
2491                    }
2492                    return format!("{prefix}: {v}");
2493                }
2494            }
2495            return v.to_string();
2496        }
2497    }
2498    String::new()
2499}
2500
2501/// One imperative line, no control characters, capped: the only shape an
2502/// authored lesson may take. A literal newline could otherwise smuggle a
2503/// second CAL statement past review (belt: serde_json escapes it anyway) or
2504/// break the one-line prompt rendering hosts assume.
2505fn sanitize_lesson(s: &str) -> String {
2506    sanitize_line(s, crate::llm::MAX_LESSON_LEN)
2507}
2508
2509/// One line, no control characters, capped. Every free-text field the model
2510/// can put into an executable proposal goes through this: a literal newline
2511/// could otherwise smuggle a second CAL statement past review (belt:
2512/// serde_json escapes it anyway), split a one-line DEFINE across the batch the
2513/// apply path iterates, or break the one-line prompt rendering hosts assume.
2514fn sanitize_line(s: &str, max: usize) -> String {
2515    let cleaned: String =
2516        s.chars().map(|c| if c.is_control() { ' ' } else { c }).collect();
2517    crate::llm::cap(cleaned.trim(), max)
2518}
2519
2520/// A relation is an identifier, not prose — it becomes a queryable predicate,
2521/// and whitespace or quotes in one would make the Fact unfindable by the very
2522/// recall that should surface it. `None` rejects the draft's `fact` proposal.
2523fn sanitize_relation(s: &str) -> Option<String> {
2524    let r = sanitize_line(s, crate::llm::MAX_RELATION_LEN);
2525    if r.is_empty()
2526        || !r
2527            .chars()
2528            .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.' | ':'))
2529    {
2530        return None;
2531    }
2532    Some(r)
2533}
2534
2535/// A model-supplied query body is placed INSIDE a `DEFINE … AS { … }` block,
2536/// so it is the one place in the vocabulary where model text becomes part of a
2537/// statement's structure rather than its content. Two belts, because the
2538/// substrate's parser strength is not something this engine gets to assume:
2539///
2540/// - **No braces.** Closing the block early is the injection shape; a saved
2541///   RECALL/ASSEMBLE body needs no braces of its own, so refusing them costs
2542///   nothing and fails closed.
2543/// - **No destructive keyword, anywhere in the body.** `cal::contains_destructive`
2544///   scans each LINE's leading keyword, which a single-line injection slips
2545///   past by construction — so scan every token here instead.
2546///
2547/// The substrate's own `validate_cal` and the saved-query read-only
2548/// verification pass still run after this; this is the layer that does not
2549/// depend on either of them being strict.
2550fn safe_definition_body(body: &str) -> bool {
2551    if body.contains('{') || body.contains('}') {
2552        return false;
2553    }
2554    !body
2555        .split(|c: char| !c.is_ascii_alphanumeric() && c != '_')
2556        .any(|tok| {
2557            ["FORGET", "PURGE", "DROP", "DEFINE"]
2558                .iter()
2559                .any(|kw| tok.eq_ignore_ascii_case(kw))
2560        })
2561}
2562
2563/// The claim GROUND entails and VERIFY stress-tests. When the draft carries a
2564/// resolvable proposal the claim names exactly what an apply would do, so what
2565/// survives the gates is what gets written — never a summary standing in for
2566/// a change the gates never saw.
2567fn claim_text(d: &crate::llm::LlmDraft, resolved: Option<&ResolvedProposal>) -> String {
2568    let summary = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
2569    match resolved {
2570        Some(r) => format!("{summary} {}", r.rendered),
2571        None => summary,
2572    }
2573}
2574
2575/// A DISCOVER draft that survived structural validation, carrying the
2576/// executable form of its proposal. Resolution happens BEFORE GROUND/VERIFY,
2577/// so both gates judge exactly what an apply would do — the rule the authored
2578/// lesson already followed, generalized to the whole vocabulary. It also means
2579/// a malformed proposal costs no model call: it dies here, not at apply.
2580struct ValidatedDraft {
2581    draft: crate::llm::LlmDraft,
2582    target_ref: String,
2583    cited: Vec<String>,
2584    resolved: Option<ResolvedProposal>,
2585}
2586
2587/// The executable shape of a validated draft. `None` on a [`ValidatedDraft`]
2588/// means advisory — the model said something a human may want to see, but
2589/// nothing the engine will ever execute.
2590struct ResolvedProposal {
2591    action: ActionKind,
2592    proposal: Proposal,
2593    /// One line naming exactly what an apply would do; folded into the claim
2594    /// both gates judge and shown in the review summary.
2595    rendered: String,
2596    summary_key: &'static str,
2597    summary_args: serde_json::Map<String, Value>,
2598    rollbackable: bool,
2599    evalset_hash: Option<String>,
2600    importance: f64,
2601    /// Set for the two Fact-writing shapes. Held rather than pre-rendered so
2602    /// the grain can carry the VERIFIER's confidence, which is not known until
2603    /// after the gates have run.
2604    fact_fields: Option<serde_json::Map<String, Value>>,
2605}
2606
2607/// The Workflow fields a `plan_revision` may touch. Thresholds and limits —
2608/// never who calls what.
2609///
2610/// The exclusion is structural, not advisory: `nodes`, `edges[].src`,
2611/// `edges[].dst` and `bindings` are simply not matchable here, so a topology
2612/// change cannot be expressed by any proposal the model can write. That keeps
2613/// a plan revision reviewable as a short list of scalar deltas rather than a
2614/// re-drawn graph, which is the difference between a reviewer checking a
2615/// number and a reviewer re-deriving a plan.
2616fn plan_edit_allowed(path: &str) -> bool {
2617    let seg: Vec<&str> = path.split('.').collect();
2618    match seg.as_slice() {
2619        ["edges", i, "cond"] | ["edges", i, "max_cycles"] => i.parse::<usize>().is_ok(),
2620        ["retries", node] => !node.is_empty(),
2621        _ => false,
2622    }
2623}
2624
2625/// Read the value at an allowlisted path (absent → `Value::Null`, which is
2626/// what an edit adding a `retries` entry must declare as its `from`).
2627fn plan_get(body: &Value, path: &str) -> Value {
2628    let mut cur = body;
2629    for seg in path.split('.') {
2630        cur = match cur {
2631            Value::Array(a) => match seg.parse::<usize>().ok().and_then(|i| a.get(i)) {
2632                Some(v) => v,
2633                None => return Value::Null,
2634            },
2635            Value::Object(o) => match o.get(seg) {
2636                Some(v) => v,
2637                None => return Value::Null,
2638            },
2639            _ => return Value::Null,
2640        };
2641    }
2642    cur.clone()
2643}
2644
2645/// Write the value at an allowlisted path. Only creates a missing key in an
2646/// object (the `retries.<node>` case); never grows an array.
2647fn plan_set(body: &mut Value, path: &str, to: Value) -> bool {
2648    let segs: Vec<&str> = path.split('.').collect();
2649    let Some((last, parents)) = segs.split_last() else {
2650        return false;
2651    };
2652    let mut cur = body;
2653    for seg in parents {
2654        cur = match cur {
2655            Value::Array(a) => match seg.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2656                Some(v) => v,
2657                None => return false,
2658            },
2659            Value::Object(o) => match o.get_mut(*seg) {
2660                Some(v) => v,
2661                None => return false,
2662            },
2663            _ => return false,
2664        };
2665    }
2666    match cur {
2667        Value::Array(a) => match last.parse::<usize>().ok().and_then(move |i| a.get_mut(i)) {
2668            Some(slot) => {
2669                *slot = to;
2670                true
2671            }
2672            None => false,
2673        },
2674        Value::Object(o) => {
2675            o.insert((*last).to_string(), to);
2676            true
2677        }
2678        _ => false,
2679    }
2680}
2681
2682/// Type-check one plan edit's new value against the field it targets. Without
2683/// this a string in `max_cycles` would be dropped by the grain deserializer
2684/// and the "applied" revision would silently mean *unlimited* — a proposal
2685/// that reads as a tightening and lands as a removal.
2686fn plan_value_ok(path: &str, to: &Value) -> bool {
2687    let seg: Vec<&str> = path.split('.').collect();
2688    match seg.as_slice() {
2689        ["edges", _, "cond"] => to.as_str().is_some_and(|c| {
2690            !c.trim().is_empty() && c.len() <= 200 && !c.chars().any(char::is_control)
2691        }),
2692        ["edges", _, "max_cycles"] | ["retries", _] => {
2693            to.as_u64().is_some_and(|n| n <= 1_000)
2694        }
2695        _ => false,
2696    }
2697}
2698
2699/// Resolve a DISCOVER draft's proposal into the executable form an apply would
2700/// run, or `None` for advisory. Every variant takes its SCOPE from the draft's
2701/// target and its supporting facts from the substrate — the model names the
2702/// change, never the subject it lands on, the namespace it lands in, or (for
2703/// code) the evalset that grades it.
2704fn resolve_proposal<S: OmsSubstrate>(
2705    sub: &S,
2706    d: &crate::llm::LlmDraft,
2707    target: &TargetRef,
2708    cited: &[String],
2709    ns_by_hash: &std::collections::BTreeMap<String, String>,
2710    caps: Capabilities,
2711    policy: &crate::policy::Policy,
2712) -> Option<ResolvedProposal> {
2713    use crate::llm::DraftProposal as P;
2714    let (skills, plans) = (&policy.skills, &policy.plans);
2715    let mut args = serde_json::Map::new();
2716    match d.parsed_proposal()? {
2717        // ---- plan: a procedure as a validated Workflow + its Skill prose ----
2718        P::Plan { description, when_to_use, nodes, edges } => {
2719            if !plans.enabled || !caps.plans {
2720                return None;
2721            }
2722            let PlanFields { skill, workflow, name, n_nodes, n_edges, existing_skill, existing_plan } =
2723                derived_plan_fields(sub, target, &description, &when_to_use, &nodes, &edges, cited, ns_by_hash, plans)?;
2724            args.insert("name".into(), Value::from(name.clone()));
2725            args.insert("nodes".into(), Value::from(n_nodes as u64));
2726            args.insert("edges".into(), Value::from(n_edges as u64));
2727            let skill_stmt = match &existing_skill {
2728                Some(h) => cal::supersede(h, "skill", &skill),
2729                None => cal::add("skill", &skill),
2730            };
2731            let plan_stmt = match &existing_plan {
2732                Some(h) => cal::supersede(h, "workflow", &workflow),
2733                None => cal::add("workflow", &workflow),
2734            };
2735            let (action, verb) = if existing_skill.is_some() || existing_plan.is_some() {
2736                (ActionKind::Revise, "revise")
2737            } else {
2738                (ActionKind::Record, "record")
2739            };
2740            Some(ResolvedProposal {
2741                action,
2742                proposal: Proposal::Cal { cal: cal::batch(&[skill_stmt, plan_stmt]) },
2743                rendered: format!(
2744                    "Proposed plan to {verb}: \"{name}\" — {n_nodes} steps, {n_edges} edges; when: {}",
2745                    skill.get("when_to_use").and_then(Value::as_str).unwrap_or("")
2746                ),
2747                summary_key: "llm.plan",
2748                summary_args: args,
2749                rollbackable: true,
2750                evalset_hash: None,
2751                importance: 0.65,
2752                fact_fields: None,
2753            })
2754        }
2755        // ---- skill: a reusable procedure from a trajectory that succeeded ----
2756        P::Skill { description, when_to_use, steps } => {
2757            if !skills.enabled {
2758                return None;
2759            }
2760            let SkillFields { fields, name, n_steps, existing } =
2761                derived_skill_fields(sub, target, &description, &when_to_use, &steps, cited, ns_by_hash, skills)?;
2762            args.insert("name".into(), Value::from(name.clone()));
2763            args.insert("steps".into(), Value::from(n_steps as u64));
2764            // A live skill of the same name in the same namespace is PATCHED
2765            // (superseded), never duplicated beside itself.
2766            let (action, cal, verb) = match existing {
2767                Some(hash) => (ActionKind::Revise, cal::supersede(&hash, "skill", &fields), "revise"),
2768                None => (ActionKind::Record, cal::add("skill", &fields), "record"),
2769            };
2770            Some(ResolvedProposal {
2771                action,
2772                proposal: Proposal::Cal { cal },
2773                rendered: format!(
2774                    "Proposed skill to {verb}: \"{name}\" — {n_steps} steps; when: {}",
2775                    fields.get("when_to_use").and_then(Value::as_str).unwrap_or("")
2776                ),
2777                summary_key: "llm.skill",
2778                summary_args: args,
2779                rollbackable: true,
2780                evalset_hash: None,
2781                importance: 0.6,
2782                fact_fields: None,
2783            })
2784        }
2785        // ---- lesson: the pre-vocabulary shape, unchanged ----
2786        P::Lesson { lesson } => {
2787            let lesson = sanitize_lesson(&lesson);
2788            if lesson.is_empty() {
2789                return None;
2790            }
2791            let fields = derived_fact_fields(target, "lesson", &lesson, cited, ns_by_hash)?;
2792            args.insert("lesson".into(), Value::from(lesson.clone()));
2793            Some(ResolvedProposal {
2794                // Same action as the deterministic lesson path — "record a
2795                // failure-derived lesson" — so dedup groups authored lessons
2796                // per target and review UIs need no new vocabulary.
2797                action: ActionKind::ClusterFailure,
2798                proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2799                rendered: format!("Proposed lesson to record: \"{lesson}\""),
2800                summary_key: "llm.lesson",
2801                summary_args: args,
2802                rollbackable: true,
2803                evalset_hash: None,
2804                importance: 0.5,
2805                fact_fields: Some(fields),
2806            })
2807        }
2808        // ---- fact: a durable fact under a model-chosen relation ----
2809        P::Fact { relation, object } => {
2810            let relation = sanitize_relation(&relation)?;
2811            let object = sanitize_line(&object, crate::llm::MAX_OBJECT_LEN);
2812            if object.is_empty() {
2813                return None;
2814            }
2815            let fields = derived_fact_fields(target, &relation, &object, cited, ns_by_hash)?;
2816            let subject = fields.get("subject").and_then(Value::as_str).unwrap_or("");
2817            args.insert("relation".into(), Value::from(relation.clone()));
2818            args.insert("object".into(), Value::from(object.clone()));
2819            Some(ResolvedProposal {
2820                action: ActionKind::Record,
2821                proposal: Proposal::Cal { cal: cal::add("fact", &fields) },
2822                rendered: format!("Proposed fact to record: {subject} {relation} \"{object}\""),
2823                summary_key: "llm.fact",
2824                summary_args: args,
2825                rollbackable: true,
2826                evalset_hash: None,
2827                importance: 0.5,
2828                fact_fields: Some(fields),
2829            })
2830        }
2831        // ---- query_revision: how the agent assembles its own context ----
2832        P::QueryRevision { body } => {
2833            let name = target.opaque();
2834            // The name comes from the target, but it still ends up inside a
2835            // statement — a quote or control character in one would change the
2836            // statement's shape rather than its content.
2837            if name.is_empty()
2838                || name.chars().any(|c| c.is_control() || c == '"' || c == '\\')
2839            {
2840                return None;
2841            }
2842            let body = sanitize_line(&body, crate::llm::MAX_QUERY_BODY_LEN);
2843            if body.is_empty() || !safe_definition_body(&body) {
2844                return None;
2845            }
2846            let stmt = match target.scheme() {
2847                "query" => format!("DEFINE QUERY \"{name}\" AS {{ {body} }}"),
2848                "template" => format!("DEFINE TEMPLATE {name} AS {{ {body} }}"),
2849                _ => return None,
2850            };
2851            // The substrate owns the grammar: if it will not parse, or will
2852            // not hand back an inverse, this is not something a reviewer
2853            // should be offered as applicable. A definition change ROLLBACK
2854            // could not undo must not be applied at all.
2855            sub.validate_cal(&stmt).ok()?;
2856            sub.definition_inverse(&stmt).ok().flatten()?;
2857            args.insert("name".into(), Value::from(name));
2858            args.insert("body".into(), Value::from(body.clone()));
2859            Some(ResolvedProposal {
2860                action: ActionKind::Revise,
2861                proposal: Proposal::Cal { cal: stmt },
2862                rendered: format!("Proposed rewrite of saved {} \"{name}\" to: {body}", target.scheme()),
2863                summary_key: "llm.query_revision",
2864                summary_args: args,
2865                rollbackable: true,
2866                evalset_hash: None,
2867                importance: 0.6,
2868                fact_fields: None,
2869            })
2870        }
2871        // ---- plan_revision: field-level edits to a Workflow grain ----
2872        P::PlanRevision { edits } => {
2873            if !caps.plans
2874                || target.scheme() != "grain"
2875                || edits.is_empty()
2876                || edits.len() > crate::llm::MAX_PLAN_EDITS
2877            {
2878                return None;
2879            }
2880            let hash = target.opaque();
2881            let g = sub.grain(hash).ok().flatten()?;
2882            if g.grain_type != "workflow" || !g.is_live() {
2883                return None;
2884            }
2885            let mut body = Value::Object(g.fields.clone());
2886            let mut deltas = Vec::new();
2887            let nodes: std::collections::BTreeSet<String> = body
2888                .get("nodes")
2889                .and_then(Value::as_array)
2890                .map(|a| a.iter().filter_map(Value::as_str).map(str::to_string).collect())
2891                .unwrap_or_default();
2892            for e in &edits {
2893                if !plan_edit_allowed(&e.path) || !plan_value_ok(&e.path, &e.to) {
2894                    return None;
2895                }
2896                // A retry count for a node that does not exist is inert, but
2897                // applying it still mints a new plan hash — and every trigger
2898                // pointing at the old one must then be walked forward. A
2899                // no-op is not worth that.
2900                if let Some(node) = e.path.strip_prefix("retries.") {
2901                    if !nodes.contains(node) {
2902                        return None;
2903                    }
2904                }
2905                // `from` is the staleness check: a proposal authored against
2906                // an older plan does not silently apply to a newer one.
2907                if plan_get(&body, &e.path) != e.from {
2908                    return None;
2909                }
2910                // A no-op edit is not a revision; it would apply, mint a new
2911                // plan hash, and orphan every trigger pointing at the old one
2912                // for nothing.
2913                if e.from == e.to {
2914                    return None;
2915                }
2916                if !plan_set(&mut body, &e.path, e.to.clone()) {
2917                    return None;
2918                }
2919                deltas.push(format!("{}: {} -> {}", e.path, e.from, e.to));
2920            }
2921            // The substrate owns the plan grammar (unique + reachable nodes,
2922            // conditions parse, every cycle bounded). An edit that would make
2923            // the plan unrunnable never reaches a reviewer as applicable.
2924            sub.validate_plan(&body).ok()?;
2925            let Value::Object(fields) = body else {
2926                return None;
2927            };
2928            let stmt = cal::supersede(hash, "workflow", &fields);
2929            // Same rule as the definition rewrite: a statement the substrate
2930            // will not accept is not something to offer a reviewer as
2931            // applicable. `validate_plan` checked the GRAPH; this checks the
2932            // statement that carries it.
2933            sub.validate_cal(&stmt).ok()?;
2934            args.insert("plan".into(), Value::from(hash));
2935            args.insert("edits".into(), Value::from(deltas.join("; ")));
2936            Some(ResolvedProposal {
2937                action: ActionKind::Revise,
2938                proposal: Proposal::Cal { cal: stmt },
2939                rendered: format!("Proposed plan revision ({})", deltas.join("; ")),
2940                summary_key: "llm.plan_revision",
2941                summary_args: args,
2942                rollbackable: true,
2943                evalset_hash: None,
2944                importance: 0.7,
2945                fact_fields: None,
2946            })
2947        }
2948        // ---- code_revision: §7.4, gated by the tool's own evalset ----
2949        P::CodeRevision { source } => {
2950            if !caps.code || target.scheme() != "tool" || source.trim().is_empty() {
2951                return None;
2952            }
2953            if source.chars().count() > crate::llm::MAX_CODE_LEN {
2954                return None;
2955            }
2956            // Rule E1's pin, resolved from the substrate. A proposer that
2957            // could name its own grader is not gated, and a tool that
2958            // declares no evalset has no gate to pass — advisory either way.
2959            let evalset = sub.tool_evalset(target.opaque()).ok().flatten()?;
2960            let mut data = serde_json::Map::new();
2961            data.insert("tool".into(), Value::from(target.opaque()));
2962            data.insert("source".into(), Value::from(source.clone()));
2963            args.insert("tool".into(), Value::from(target.opaque()));
2964            args.insert("bytes".into(), Value::from(source.len() as u64));
2965            Some(ResolvedProposal {
2966                action: ActionKind::CodeRevision,
2967                proposal: Proposal::Data { data },
2968                rendered: format!(
2969                    "Proposed new source for tool {} ({} bytes), gated by evalset {}",
2970                    target.opaque(),
2971                    source.len(),
2972                    evalset
2973                ),
2974                summary_key: "llm.code_revision",
2975                summary_args: args,
2976                rollbackable: true,
2977                evalset_hash: Some(evalset),
2978                importance: 0.8,
2979                fact_fields: None,
2980            })
2981        }
2982    }
2983}
2984
2985/// Stamp a validated DISCOVER draft as an `origin = llm` recommendation.
2986/// Default shape: an advisory `Flag` carrying `Proposal::Data` (no executable
2987/// mutation). A draft whose proposal RESOLVED (see [`resolve_proposal`])
2988/// instead stamps as that executable change, reviewable with the exact line an
2989/// apply would run. Either way `Origin::Llm` plus the no-manifest analyzer id
2990/// leave it structurally ineligible for auto-apply — and independently, no
2991/// class this vocabulary can reach except `memory` is auto-appliable at all
2992/// (`Policy::grants_auto_apply`). The only path into the agent is a human
2993/// review with a BECAUSE followed by an explicit apply.
2994fn stamp_llm(
2995    model: &str,
2996    d: &crate::llm::LlmDraft,
2997    target_ref: String,
2998    cited: Vec<String>,
2999    resolved: Option<ResolvedProposal>,
3000    confidence: f64,
3001    now_ms: i64,
3002) -> Recommendation {
3003    let summary_text = crate::llm::cap(&d.summary, crate::llm::MAX_SUMMARY_LEN);
3004    let guidance = if d.guidance.trim().is_empty() {
3005        None
3006    } else {
3007        Some(crate::llm::cap(&d.guidance, crate::llm::MAX_GUIDANCE_LEN))
3008    };
3009    let (action, proposal, summary, rollbackable, importance, evalset_hash, content) = match resolved {
3010        Some(mut r) => {
3011            // What the proposal would DO, before the verifier's confidence
3012            // is folded into the fact: the dedup key fingerprints this, so
3013            // the same lesson at a different confidence is one finding.
3014            let content = match &r.fact_fields {
3015                Some(fields) => format!(
3016                    "{} {}",
3017                    fields.get("relation").and_then(Value::as_str).unwrap_or(""),
3018                    fields.get("object").and_then(Value::as_str).unwrap_or("")
3019                ),
3020                None => match &r.proposal {
3021                    Proposal::Cal { cal } => cal.clone(),
3022                    Proposal::Data { data } => Value::Object(data.clone()).to_string(),
3023                    Proposal::Edit { diff, .. } => diff.clone(),
3024                },
3025            };
3026            // The grain records the VERIFIER's calibrated confidence — the
3027            // independent signal — never the proposer's self-report.
3028            if let Some(mut fields) = r.fact_fields.take() {
3029                fields.insert("confidence".into(), Value::from(confidence.clamp(0.0, 1.0)));
3030                r.proposal = Proposal::Cal { cal: cal::add("fact", &fields) };
3031            }
3032            let mut args = r.summary_args;
3033            args.insert("text".into(), Value::from(summary_text));
3034            (
3035                r.action,
3036                r.proposal,
3037                Summary::new(r.summary_key, args),
3038                r.rollbackable,
3039                r.importance,
3040                r.evalset_hash,
3041                Some(content),
3042            )
3043        }
3044        None => {
3045            let mut args = serde_json::Map::new();
3046            args.insert("text".into(), Value::from(summary_text));
3047            let mut data = serde_json::Map::new();
3048            data.insert("source".into(), Value::from("llm"));
3049            (
3050                ActionKind::Flag,
3051                Proposal::Data { data },
3052                Summary::new("llm.discover", args),
3053                false,
3054                0.3,
3055                None,
3056                None,
3057            )
3058        }
3059    };
3060    // An advisory flag keeps the analyzer-style key (one open flag per
3061    // target); an executable proposal keys on its content too, because
3062    // there the content is the finding.
3063    let dedup = match &content {
3064        Some(c) => crate::recommendation::authored_dedup_key("llm", &target_ref, action, c),
3065        None => dedup_key("llm", &target_ref, action),
3066    };
3067    Recommendation {
3068        hash: String::new(),
3069        analyzer: "loop.llm/1".to_string(),
3070        params_snapshot: serde_json::Map::new(),
3071        origin: Origin::Llm { model: model.to_string() },
3072        target_ref: target_ref.clone(),
3073        action_kind: action,
3074        dedup_key: dedup,
3075        summary,
3076        severity: Severity::Low,
3077        proposal,
3078        destructive: false,
3079        rollbackable,
3080        evidence: cited,
3081        evidence_query: None,
3082        metric: None,
3083        // The verifier's calibrated confidence — not a hardcoded default.
3084        confidence: confidence.clamp(0.0, 1.0),
3085        importance,
3086        created_at_ms: now_ms,
3087        guidance,
3088        evalset_hash,
3089        status: RecStatus::Pending,
3090    }
3091}
3092
3093/// The Fact an approved `lesson` or `fact` proposal writes: subject from the
3094/// entity target, the model's relation (`"lesson"` for the lesson shape —
3095/// prescriptive prose, distinct from the deterministic `fails_with` signature
3096/// facts), the sanitized object, and the DOMINANT namespace of the cited
3097/// evidence (max count, ties to the lexicographically smallest — the
3098/// tool_failure rule), never a namespace the model names. `None` when the
3099/// target gives no subject.
3100/// The `skill` paragraph appended to the DISCOVER instructions when the host
3101/// allows skill authoring. Kept beside the fixed instruction text it extends.
3102fn skill_instructions(min_steps: u32) -> String {
3103    format!(
3104        " (6) {{\"kind\":\"skill\",\"description\":\"...\",\"when_to_use\":\"...\",\
3105\"steps\":[\"...\",\"...\"]}} with target \"entity:<ns>/<skill-name>\" — a REUSABLE \
3106PROCEDURE the agent carried out successfully in the evidence: a sequence of tool \
3107calls that reached its goal, which a later session facing the same situation \
3108should not have to rediscover. Give {min_steps} to {} ordered steps, each naming \
3109the tool called and the values that mattered (the field checked, the tag set, the \
3110exact format produced), a one-line description, and 'when_to_use' — the situation \
3111that should trigger it. The skill-name is a short identifier (letters, digits, \
3112_ -). If a saved skill already covers this procedure, use ITS name so it is \
3113patched rather than duplicated. Do not propose a skill for a procedure that \
3114failed, or for one already saved and unchanged.",
3115        crate::llm::MAX_SKILL_STEPS
3116    )
3117}
3118
3119/// The `plan` paragraph appended to the DISCOVER instructions when the host
3120/// allows plan authoring. The condition grammar is the runtime's frozen v1
3121/// grammar, stated so the model writes conditions the plan validator accepts.
3122fn plan_instructions(min_nodes: u32) -> String {
3123    format!(
3124        " (7) {{\"kind\":\"plan\",\"description\":\"...\",\"when_to_use\":\"...\",\
3125\"nodes\":[{{\"id\":\"list_open\",\"tool\":\"<tool name>\",\"step\":\"...\"}},...],\
3126\"edges\":[{{\"src\":\"list_open\",\"dst\":\"tag\",\"cond\":\"shared_incident == true\"}},...]}} \
3127with target \"entity:<ns>/<plan-name>\" — the same reusable procedure as a skill, \
3128but as a PLAN the runtime can validate and run: {min_nodes} to {} steps, each an \
3129'id' (letters, digits, _ -), the 'tool' it calls — which MUST be a tool named in \
3130the cited evidence — and what the step does with it; and 'edges' from step to \
3131step. An edge 'cond' is optional and uses exactly this grammar: 'path == literal', \
3132'path != literal', 'path exists' or '!path', where path is dotted names and the \
3133literal is a JSON string, number, true, false or null — no other operators; state \
3134a threshold as a flag the step sets ('reporters_ge_3 == true'). A loop back to an \
3135earlier step needs 'max_cycles'. Prefer a plan over a skill when the procedure \
3136has branches or a loop; prefer a skill when it is a straight list. If a saved \
3137plan already covers this procedure, use ITS name so it is patched.",
3138        crate::llm::MAX_PLAN_NODES
3139    )
3140}
3141
3142/// What a `plan` proposal resolves to: the Skill grain (prose), the Workflow
3143/// grain (structure), the name both take from the target, the counts, and the
3144/// live pair of that name (to supersede) if there is one.
3145struct PlanFields {
3146    skill: serde_json::Map<String, Value>,
3147    workflow: serde_json::Map<String, Value>,
3148    name: String,
3149    n_nodes: usize,
3150    n_edges: usize,
3151    existing_skill: Option<String>,
3152    existing_plan: Option<String>,
3153}
3154
3155/// The two grains a `plan` proposal would write.
3156///
3157/// Grounding is structural: every step's tool must be one the cited evidence
3158/// shows was called, so the model cannot plan around a tool it invented. The
3159/// workflow body is handed to the substrate's own plan validator before the
3160/// draft can be stamped applicable — a plan a reviewer could approve is one
3161/// the runtime would accept. The pair shares a `name`; the Workflow carries it
3162/// as a host field (the type is a container by design), and that is how the
3163/// live plan of a name is found to be patched rather than duplicated.
3164#[allow(clippy::too_many_arguments)]
3165fn derived_plan_fields<S: SubstrateRead>(
3166    sub: &S,
3167    target: &TargetRef,
3168    description: &str,
3169    when_to_use: &str,
3170    nodes: &[crate::llm::PlanNodeDraft],
3171    edges: &[crate::llm::PlanEdgeDraft],
3172    cited: &[String],
3173    ns_by_hash: &std::collections::BTreeMap<String, String>,
3174    plans: &crate::policy::PlanAuthoring,
3175) -> Option<PlanFields> {
3176    if target.scheme() != "entity" {
3177        return None;
3178    }
3179    let name = sanitize_skill_name(
3180        target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3181    )?;
3182    let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3183    let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3184    if description.is_empty() || when_to_use.is_empty() {
3185        return None;
3186    }
3187    if nodes.len() < plans.min_nodes.max(1) as usize || nodes.len() > crate::llm::MAX_PLAN_NODES {
3188        return None;
3189    }
3190    // The tools the evidence shows were actually called.
3191    let known_tools: BTreeSet<String> = cited
3192        .iter()
3193        .filter_map(|h| sub.grain(h).ok().flatten())
3194        .filter_map(|g| g.tool_name().map(normalize_ident))
3195        .collect();
3196    let mut ids: Vec<String> = Vec::new();
3197    let mut steps: Vec<String> = Vec::new();
3198    let mut seen: BTreeSet<String> = BTreeSet::new();
3199    for n in nodes {
3200        let id = sanitize_skill_name(&n.id)?;
3201        if !seen.insert(id.clone()) {
3202            return None; // duplicate step id
3203        }
3204        let tool = sanitize_line(&n.tool, crate::llm::MAX_SKILL_NAME_LEN);
3205        if tool.is_empty() || !known_tools.contains(&normalize_ident(&tool)) {
3206            return None; // a tool the evidence never shows: not grounded
3207        }
3208        let step = sanitize_line(&n.step, crate::llm::MAX_SKILL_STEP_LEN);
3209        if step.is_empty() {
3210            return None;
3211        }
3212        steps.push(format!("{}. {id} [{tool}]: {step}", steps.len() + 1));
3213        ids.push(id);
3214    }
3215    let mut edge_vals: Vec<Value> = Vec::new();
3216    let mut flow_lines: Vec<String> = Vec::new();
3217    for e in edges {
3218        let src = sanitize_skill_name(&e.src)?;
3219        let dst = sanitize_skill_name(&e.dst)?;
3220        if !seen.contains(&src) || !seen.contains(&dst) {
3221            return None;
3222        }
3223        let mut ev = serde_json::Map::new();
3224        ev.insert("src".into(), Value::from(src.clone()));
3225        ev.insert("dst".into(), Value::from(dst.clone()));
3226        let mut label = format!("{src} → {dst}");
3227        if let Some(c) = e
3228            .cond
3229            .as_deref()
3230            .map(|c| sanitize_line(c, crate::llm::MAX_COND_LEN))
3231            .filter(|c| !c.is_empty())
3232        {
3233            label.push_str(&format!(" if {c}"));
3234            ev.insert("cond".into(), Value::from(c));
3235        }
3236        if let Some(m) = e.max_cycles {
3237            if m == 0 || m > 100 {
3238                return None;
3239            }
3240            label.push_str(&format!(" (at most {m} times)"));
3241            ev.insert("max_cycles".into(), Value::from(m));
3242        }
3243        flow_lines.push(label);
3244        edge_vals.push(Value::Object(ev));
3245    }
3246    if edge_vals.len() > 4 * ids.len() {
3247        return None;
3248    }
3249    // Namespace: where the evidence lives, by majority — a lesson's rule.
3250    let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3251    for h in cited {
3252        if let Some(ns) = ns_by_hash.get(h) {
3253            if !ns.is_empty() {
3254                *ns_counts.entry(ns.as_str()).or_default() += 1;
3255            }
3256        }
3257    }
3258    let ns = ns_counts
3259        .iter()
3260        .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3261        .map(|(ns, _)| ns.to_string());
3262
3263    // The Workflow: what the runtime validates. Unbound steps are abstract
3264    // nodes — legal, and what a plan over a host's own tools is.
3265    let mut workflow = serde_json::Map::new();
3266    workflow.insert("nodes".into(), Value::from(ids.clone()));
3267    workflow.insert("edges".into(), Value::Array(edge_vals));
3268    workflow.insert("name".into(), Value::from(name.clone()));
3269    if let Some(ns) = &ns {
3270        workflow.insert("namespace".into(), Value::from(ns.clone()));
3271    }
3272    sub.validate_plan(&Value::Object(workflow.clone())).ok()?;
3273
3274    // The Skill: the same procedure as prose, with the graph's edges spelled
3275    // out under the steps so a reader sees the branches the plan encodes.
3276    let mut instructions = steps.join("\n");
3277    if !flow_lines.is_empty() {
3278        instructions.push_str("\n\nFlow:\n");
3279        instructions.push_str(&flow_lines.iter().map(|l| format!("- {l}")).collect::<Vec<_>>().join("\n"));
3280    }
3281    let mut skill = serde_json::Map::new();
3282    skill.insert("name".into(), Value::from(name.clone()));
3283    skill.insert("description".into(), Value::from(description));
3284    skill.insert("when_to_use".into(), Value::from(when_to_use));
3285    skill.insert("instructions".into(), Value::from(instructions));
3286    if let Some(ns) = &ns {
3287        skill.insert("namespace".into(), Value::from(ns.clone()));
3288    }
3289    let live = |gt: &str, pick: &dyn Fn(&GrainRecord) -> bool| -> Option<String> {
3290        sub.grains_of_type(gt, ns.as_deref(), ReadOpts { live_only: true, since_ms: None })
3291            .ok()?
3292            .into_iter()
3293            .find(|g| pick(g))
3294            .map(|g| g.hash)
3295    };
3296    let existing_skill = live(crate::model::grain_type::SKILL, &|g| g.skill_name() == Some(name.as_str()));
3297    let existing_plan = live(crate::model::grain_type::WORKFLOW, &|g| g.str_field("name") == Some(name.as_str()));
3298    Some(PlanFields {
3299        skill,
3300        workflow,
3301        name,
3302        n_nodes: ids.len(),
3303        n_edges: flow_lines.len(),
3304        existing_skill,
3305        existing_plan,
3306    })
3307}
3308
3309/// The metric name under which the Verify gate records a premise that moved.
3310pub const PREMISE_DRIFT_METRIC: &str = "premise_drift";
3311
3312/// The Verify gate's second question. For every applied recommendation, the
3313/// grains it cited are looked up again: one that has been retracted, or
3314/// superseded by a grain holding a DIFFERENT value, is a premise that moved.
3315/// A value-identical supersession — what consolidation does — is not, and
3316/// neither is a supersession the recommendation's OWN apply performed: a
3317/// contradiction resolution cites the two conflicting facts and retires one
3318/// of them; that is the change it was approved to make, not its premise
3319/// moving out from under it.
3320///
3321/// Compared with the wrong rule (counting the apply's own work as drift), this
3322/// is what keeps the gate quiet on the analyzers that exist to supersede.
3323/// Records `drifted` in the outcome series once per distinct count (so a
3324/// pass does not re-record what the last pass already did) and returns an
3325/// input `outcome_review` turns into the revert proposal. The reviewer
3326/// decides; nothing here applies.
3327fn detect_premise_drift<S: OmsSubstrate>(
3328    sub: &S,
3329    p: &mut LoopPersisted,
3330    now_ms: i64,
3331) -> Result<Vec<OutcomeInput>> {
3332    let mut out = Vec::new();
3333    let applied: Vec<(String, String, Vec<String>)> = p
3334        .applied
3335        .iter()
3336        .filter(|(h, _)| p.status_index.get(*h) == Some(&RecStatus::Applied))
3337        .map(|(h, a)| (h.clone(), a.target_ref.clone(), a.created_hashes.clone()))
3338        .collect();
3339    for (rec_hash, target_ref, own) in applied {
3340        let Ok(rec) = load_rec(sub, &rec_hash) else { continue };
3341        if rec.evidence.is_empty() {
3342            continue;
3343        }
3344        let mut moved = 0u64;
3345        for e in &rec.evidence {
3346            match sub.grain(e)? {
3347                None => moved += 1, // retracted or gone
3348                Some(g) => {
3349                    let Some(newer) = &g.superseded_by else { continue };
3350                    if own.iter().any(|c| c == newer) {
3351                        continue; // the apply's own supersession
3352                    }
3353                    match sub.grain(newer)? {
3354                        // Superseded by something unreadable: a retraction
3355                        // (the reference substrate marks FORGET this way).
3356                        None => moved += 1,
3357                        Some(n) => {
3358                            if !same_value(&g, &n) {
3359                                moved += 1;
3360                            }
3361                        }
3362                    }
3363                }
3364            }
3365        }
3366        if moved == 0 {
3367            continue;
3368        }
3369        let already = p
3370            .outcomes
3371            .get(&rec_hash)
3372            .and_then(|v| v.iter().rev().find(|o| o.metric == PREMISE_DRIFT_METRIC))
3373            .is_some_and(|o| o.current == moved as f64);
3374        if !already {
3375            p.outcomes.entry(rec_hash.clone()).or_default().push(
3376                crate::recommendation::OutcomeResult {
3377                    rec_hash: rec_hash.clone(),
3378                    metric: PREMISE_DRIFT_METRIC.into(),
3379                    baseline: 0.0,
3380                    current: moved as f64,
3381                    verdict: "drifted".into(),
3382                    horizon_ms: 0,
3383                    checkpoint: None,
3384                    measured_at_ms: now_ms,
3385                },
3386            );
3387        }
3388        out.push(OutcomeInput {
3389            rec_hash,
3390            target_ref,
3391            metric: PREMISE_DRIFT_METRIC.into(),
3392            baseline: 0.0,
3393            current: moved as f64,
3394            unit: "superseded premises".into(),
3395            higher_is_better: false,
3396        });
3397    }
3398    Ok(out)
3399}
3400
3401/// Does the superseding grain say the same thing as the one it replaced? A
3402/// fact compares its object; anything else compares its text body. Two grains
3403/// that cannot be compared are treated as different — the fail-closed
3404/// reading, since a premise we cannot confirm still holds is one that moved.
3405fn same_value(old: &GrainRecord, new: &GrainRecord) -> bool {
3406    if let (Some(a), Some(b)) = (old.fact_object(), new.fact_object()) {
3407        return normalize_ident(a) == normalize_ident(b);
3408    }
3409    for key in ["content", "tool_content", "body", "text", "object"] {
3410        if let (Some(a), Some(b)) = (old.str_field(key), new.str_field(key)) {
3411            return normalize_ident(a) == normalize_ident(b);
3412        }
3413    }
3414    false
3415}
3416
3417/// A skill name is an identifier: `[A-Za-z0-9_-]`, bounded, case preserved.
3418fn sanitize_skill_name(s: &str) -> Option<String> {
3419    let t = s.trim();
3420    if t.is_empty()
3421        || t.chars().count() > crate::llm::MAX_SKILL_NAME_LEN
3422        || !t.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
3423    {
3424        return None;
3425    }
3426    Some(t.to_string())
3427}
3428
3429/// What a `skill` proposal resolves to: the grain's fields, the name it took
3430/// from the target, how many steps survived sanitizing, and the live skill of
3431/// that name in the same namespace (to supersede) if there is one.
3432struct SkillFields {
3433    fields: serde_json::Map<String, Value>,
3434    name: String,
3435    n_steps: usize,
3436    existing: Option<String>,
3437}
3438
3439/// The fields of the Skill grain a `skill` proposal would write.
3440///
3441/// The namespace is the one most of the cited evidence lives in — the same
3442/// rule a lesson follows — never one the model names.
3443#[allow(clippy::too_many_arguments)]
3444fn derived_skill_fields<S: SubstrateRead>(
3445    sub: &S,
3446    target: &TargetRef,
3447    description: &str,
3448    when_to_use: &str,
3449    steps: &[String],
3450    cited: &[String],
3451    ns_by_hash: &std::collections::BTreeMap<String, String>,
3452    skills: &crate::policy::SkillAuthoring,
3453) -> Option<SkillFields> {
3454    if target.scheme() != "entity" {
3455        return None;
3456    }
3457    let name = sanitize_skill_name(
3458        target.opaque().rsplit_once('/').map(|(_, n)| n).unwrap_or(target.opaque()),
3459    )?;
3460    let description = sanitize_line(description, crate::llm::MAX_OBJECT_LEN);
3461    let when_to_use = sanitize_line(when_to_use, crate::llm::MAX_OBJECT_LEN);
3462    let steps: Vec<String> = steps
3463        .iter()
3464        .map(|st| sanitize_line(st, crate::llm::MAX_SKILL_STEP_LEN))
3465        .filter(|st| !st.is_empty())
3466        .take(crate::llm::MAX_SKILL_STEPS)
3467        .collect();
3468    if description.is_empty() || when_to_use.is_empty() || steps.len() < skills.min_steps.max(1) as usize {
3469        return None;
3470    }
3471    // Namespace: where the evidence lives, by majority (ties → lexically
3472    // first), exactly as a lesson's.
3473    let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3474    for h in cited {
3475        if let Some(ns) = ns_by_hash.get(h) {
3476            if !ns.is_empty() {
3477                *ns_counts.entry(ns.as_str()).or_default() += 1;
3478            }
3479        }
3480    }
3481    let ns = ns_counts
3482        .iter()
3483        .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3484        .map(|(ns, _)| ns.to_string());
3485    let instructions = steps
3486        .iter()
3487        .enumerate()
3488        .map(|(i, st)| format!("{}. {st}", i + 1))
3489        .collect::<Vec<_>>()
3490        .join("\n");
3491    let mut fields = serde_json::Map::new();
3492    fields.insert("name".into(), Value::from(name.clone()));
3493    fields.insert("description".into(), Value::from(description));
3494    fields.insert("when_to_use".into(), Value::from(when_to_use));
3495    fields.insert("instructions".into(), Value::from(instructions));
3496    if let Some(ns) = &ns {
3497        fields.insert("namespace".into(), Value::from(ns.clone()));
3498    }
3499    // Patch, don't duplicate: the live skill of this name in this namespace.
3500    let existing = sub
3501        .grains_of_type(
3502            crate::model::grain_type::SKILL,
3503            ns.as_deref(),
3504            ReadOpts { live_only: true, since_ms: None },
3505        )
3506        .ok()?
3507        .into_iter()
3508        .find(|g| g.skill_name() == Some(name.as_str()))
3509        .map(|g| g.hash);
3510    Some(SkillFields { fields, name, n_steps: steps.len(), existing })
3511}
3512
3513fn derived_fact_fields(
3514    target: &TargetRef,
3515    relation: &str,
3516    object: &str,
3517    cited: &[String],
3518    ns_by_hash: &std::collections::BTreeMap<String, String>,
3519) -> Option<serde_json::Map<String, Value>> {
3520    if target.scheme() != "entity" {
3521        return None;
3522    }
3523    let subject = target
3524        .opaque()
3525        .rsplit_once('/')
3526        .map(|(_, s)| s)
3527        .unwrap_or(target.opaque());
3528    if subject.is_empty() {
3529        return None;
3530    }
3531    let mut ns_counts: std::collections::BTreeMap<&str, usize> = Default::default();
3532    for h in cited {
3533        if let Some(ns) = ns_by_hash.get(h) {
3534            if !ns.is_empty() {
3535                *ns_counts.entry(ns.as_str()).or_default() += 1;
3536            }
3537        }
3538    }
3539    let lesson_ns = ns_counts
3540        .iter()
3541        .max_by(|a, b| a.1.cmp(b.1).then_with(|| b.0.cmp(a.0)))
3542        .map(|(ns, _)| ns.to_string());
3543    let mut fields = serde_json::Map::new();
3544    fields.insert("subject".into(), Value::from(subject));
3545    fields.insert("relation".into(), Value::from(relation));
3546    fields.insert("object".into(), Value::from(object));
3547    // `confidence` is stamped by the caller from the VERIFIER's calibrated
3548    // score, not the proposer's self-report — so it is deliberately absent
3549    // here, where only the proposer has spoken.
3550    if let Some(ns) = lesson_ns {
3551        fields.insert("namespace".into(), Value::from(ns));
3552    }
3553    Some(fields)
3554}
3555
3556/// The action kinds that apply ONLY through the evalset-run gating edge
3557/// (§7.4 for tool code; the tuning seam's adapter promotion inherits the
3558/// same rule). One predicate so the gate, the rollbackable stamp, and the
3559/// promotion write can never disagree on membership.
3560fn requires_gating(kind: ActionKind) -> bool {
3561    matches!(
3562        kind,
3563        ActionKind::CodeRevision | ActionKind::AdapterRevision
3564    )
3565}
3566
3567fn stamp(
3568    m: &AnalyzerManifest,
3569    params: &crate::manifest::Params,
3570    d: crate::recommendation::RecDraft,
3571    now_ms: i64,
3572) -> Result<Recommendation> {
3573    let target = TargetRef::parse(&d.target_ref)?;
3574    // Rule E1 at the door (§7.4): an unpinned code revision, a pinned
3575    // non-code target, or a code action on a non-tool target never becomes
3576    // a recommendation at all.
3577    crate::recommendation::validate_code_rules(
3578        d.action_kind,
3579        target.target_class(),
3580        d.evalset_hash.as_deref(),
3581    )?;
3582    // A revert's identity is the recommendation it retracts, not just its
3583    // target: two regressed lessons on one entity are two reverts.
3584    let revert_of = match (&d.action_kind, &d.proposal) {
3585        (ActionKind::Revert, Proposal::Data { data }) => {
3586            data.get("revert_of").and_then(|v| v.as_str()).map(str::to_string)
3587        }
3588        _ => None,
3589    };
3590    let dedup = match revert_of.as_deref() {
3591        Some(h) => crate::recommendation::revert_dedup_key(m.family(), &d.target_ref, h),
3592        None => dedup_key(m.family(), &d.target_ref, d.action_kind),
3593    };
3594    let destructive = match &d.proposal {
3595        Proposal::Cal { cal } => cal::contains_destructive(cal),
3596        _ => false,
3597    };
3598    let rollbackable = match &d.proposal {
3599        Proposal::Cal { .. } => !destructive,
3600        Proposal::Edit { .. } => false,
3601        // A code or adapter revision applies by WRITING the promotion
3602        // grain; retracting it is the exact inverse — rollbackable by
3603        // construction.
3604        Proposal::Data { .. } => requires_gating(d.action_kind),
3605    };
3606    let mut evidence = d.evidence;
3607    evidence.truncate(MAX_EVIDENCE);
3608    // Provenance follows the analyzer's trust class: a subprocess
3609    // (`--analyzer-cmd`) finding is stamped `Command` — structurally
3610    // auto-apply-ineligible and badged [external] on the recall surface —
3611    // not `Builtin`.
3612    let origin = match m.trust_class {
3613        crate::manifest::TrustClass::Command => Origin::Command { id: m.id.clone() },
3614        _ => Origin::Builtin,
3615    };
3616    Ok(Recommendation {
3617        hash: String::new(),
3618        analyzer: m.id.clone(),
3619        params_snapshot: params.snapshot(),
3620        origin,
3621        target_ref: target.as_string(),
3622        action_kind: d.action_kind,
3623        dedup_key: dedup,
3624        summary: d.summary,
3625        severity: d.severity,
3626        proposal: d.proposal,
3627        destructive,
3628        rollbackable,
3629        evidence,
3630        evidence_query: d.evidence_query,
3631        metric: d.metric,
3632        confidence: d.confidence,
3633        importance: d.importance,
3634        created_at_ms: now_ms,
3635        guidance: None,
3636        evalset_hash: d.evalset_hash,
3637        status: RecStatus::Pending,
3638    })
3639}
3640
3641fn validate_because(because: &str) -> Result<String> {
3642    let trimmed = because.trim();
3643    if trimmed.is_empty() {
3644        return Err(Error::InvalidProposal(
3645            "a BECAUSE reason is required".into(),
3646        ));
3647    }
3648    if trimmed.chars().count() > MAX_BECAUSE {
3649        return Err(Error::InvalidProposal(format!(
3650            "BECAUSE exceeds {MAX_BECAUSE} chars"
3651        )));
3652    }
3653    Ok(trimmed.to_string())
3654}
3655
3656fn missing_capability(m: &AnalyzerManifest, caps: Capabilities) -> Option<&'static str> {
3657    for req in &m.requires {
3658        match req {
3659            Capability::Forks if !caps.forks => return Some("forks"),
3660            Capability::Telemetry if !caps.telemetry => return Some("telemetry"),
3661            Capability::Embeddings if !caps.embeddings => return Some("embeddings"),
3662            _ => {}
3663        }
3664    }
3665    None
3666}
3667
3668fn severity_floor_for(p: &LoopPersisted, analyzer_id: &str) -> Option<Severity> {
3669    p.config.get(analyzer_id).and_then(|c| c.severity_floor)
3670}
3671
3672fn gate(
3673    opts: &RunOptions,
3674    p: &LoopPersisted,
3675    new_grains: u64,
3676    new_errors: u64,
3677    now_ms: i64,
3678) -> Option<SkipReason> {
3679    let any = opts.min_new.is_some() || opts.min_new_errors.is_some() || opts.if_stale_ms.is_some();
3680    if !any {
3681        return None;
3682    }
3683    let min_new_ok = opts.min_new.is_some_and(|m| new_grains >= m);
3684    let min_err_ok = opts.min_new_errors.is_some_and(|m| new_errors >= m);
3685    let stale_ok = opts
3686        .if_stale_ms
3687        .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
3688    if min_new_ok || min_err_ok || stale_ok {
3689        return None;
3690    }
3691    // Choose the honest reason: staleness-only gate → not_stale, else min_new.
3692    if opts.if_stale_ms.is_some() && opts.min_new.is_none() && opts.min_new_errors.is_none() {
3693        Some(SkipReason::NotStale)
3694    } else {
3695        Some(SkipReason::MinNewNotMet)
3696    }
3697}
3698
3699/// What landed since a watermark, in every unit a gate can count.
3700#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
3701pub(crate) struct NewSince {
3702    /// Grains of the four evidence-bearing types.
3703    pub grains: u64,
3704    /// Tool grains recording a failure (`--min-new-errors`).
3705    pub error_events: u64,
3706    /// Event grains — turns, in a chat deployment (`cadence.every_events`).
3707    pub events: u64,
3708    /// Distinct `session_id`s among those Events (`cadence.every_sessions`).
3709    pub sessions: u64,
3710}
3711
3712fn count_new<S: SubstrateRead>(sub: &S, watermark: Option<i64>) -> Result<NewSince> {
3713    let opts = ReadOpts {
3714        live_only: false,
3715        since_ms: watermark.map(|w| w + 1),
3716    };
3717    let mut n = NewSince::default();
3718    let mut sessions: BTreeSet<&str> = BTreeSet::new();
3719    let mut events_held: Vec<GrainRecord> = Vec::new();
3720    for t in [
3721        crate::model::grain_type::FACT,
3722        crate::model::grain_type::EVENT,
3723        crate::model::grain_type::TOOL,
3724        crate::model::grain_type::OBSERVATION,
3725    ] {
3726        let g = sub.grains_of_type(t, None, opts)?;
3727        n.grains += g.len() as u64;
3728        // The error gate (--min-new-errors) watches captured tool failures.
3729        if t == crate::model::grain_type::TOOL {
3730            n.error_events += g.iter().filter(|e| e.is_error()).count() as u64;
3731        }
3732        if t == crate::model::grain_type::EVENT {
3733            n.events = g.len() as u64;
3734            events_held = g;
3735        }
3736    }
3737    for e in &events_held {
3738        if let Some(sid) = e.str_field("session_id").filter(|s| !s.is_empty()) {
3739            sessions.insert(sid);
3740        }
3741    }
3742    n.sessions = sessions.len() as u64;
3743    Ok(n)
3744}
3745
3746/// The policy cadence, evaluated: `None` when a pass is due.
3747///
3748/// OR over the thresholds the host set — the pass runs when any one is met.
3749/// The same shape as the per-call gate above, deliberately: the flags and
3750/// the policy block are one mechanism spelled in two places, and the flags
3751/// win when both are present.
3752fn cadence_gate(c: &crate::policy::Cadence, p: &LoopPersisted, new: NewSince, now_ms: i64) -> Option<SkipReason> {
3753    if !c.is_set() {
3754        return None;
3755    }
3756    let time_ok = c
3757        .every_ms
3758        .is_some_and(|d| p.state.last_run_ms.is_none_or(|last| now_ms - last >= d));
3759    let grains_ok = c.every_grains.is_some_and(|m| new.grains >= m);
3760    let events_ok = c.every_events.is_some_and(|m| new.events >= m);
3761    let sessions_ok = c.every_sessions.is_some_and(|m| new.sessions >= m);
3762    if time_ok || grains_ok || events_ok || sessions_ok {
3763        None
3764    } else {
3765        Some(SkipReason::CadenceNotDue)
3766    }
3767}
3768
3769fn existing_dedup_keys<S: SubstrateRead>(sub: &S, p: &LoopPersisted) -> Result<BTreeSet<String>> {
3770    let grains = sub.grains_of_type(
3771        crate::model::grain_type::RECOMMENDATION,
3772        Some(LOOP_NS),
3773        ReadOpts {
3774            live_only: false,
3775            since_ms: None,
3776        },
3777    )?;
3778    let mut set = BTreeSet::new();
3779    for g in grains {
3780        let status = p
3781            .status_index
3782            .get(&g.hash)
3783            .copied()
3784            .unwrap_or(RecStatus::Pending);
3785        // Pending/approved (still open) and applied (already handled)
3786        // recommendations suppress re-proposal of the same finding. Rejected
3787        // is handled by cooldowns, and so is a rollback the Verify gate
3788        // caused (`strike_cooldown` at the revert apply); an operator's own
3789        // rollback and expiry may legitimately re-propose (the situation
3790        // returned).
3791        if matches!(
3792            status,
3793            RecStatus::Pending | RecStatus::Approved | RecStatus::Applied
3794        ) {
3795            if let Some(key) = g.str_field("dedup_key") {
3796                set.insert(key.to_string());
3797            }
3798        }
3799    }
3800    Ok(set)
3801}
3802
3803/// Put a finding's `dedup_key` on an exponential cooldown: 7d, 14d, 28d, …
3804/// capped at 90d, so a finding a reviewer keeps rejecting stops re-surfacing
3805/// on a fixed 7d cadence (it was a flat 7d despite the "doubling" comment).
3806/// Two events earn a strike: a reviewer's rejection, and a revert the Verify
3807/// gate proposed on a measured regression — both are a verdict that the
3808/// finding, as it stands, should not come back on the next pass.
3809fn strike_cooldown(p: &mut LoopPersisted, dedup_key: String, now_ms: i64) {
3810    const BASE_MS: i64 = 7 * 86_400_000;
3811    const CAP_MS: i64 = 90 * 86_400_000;
3812    let strikes = p.cooldown_strikes.entry(dedup_key.clone()).or_insert(0);
3813    let interval = BASE_MS.saturating_mul(1_i64 << (*strikes).min(31)).min(CAP_MS);
3814    *strikes = strikes.saturating_add(1);
3815    p.cooldowns.insert(dedup_key, now_ms + interval);
3816}
3817
3818fn load_rec<S: SubstrateRead>(sub: &S, rec_hash: &str) -> Result<Recommendation> {
3819    let g = sub
3820        .grain(rec_hash)?
3821        .ok_or_else(|| Error::NotFound(rec_hash.into()))?;
3822    Recommendation::from_fields(rec_hash, &g.fields)
3823}
3824
3825/// Is this line a definition rewrite — a statement that changes a saved
3826/// `qry:`/`tpl:` registry row rather than writing a grain?
3827///
3828/// A keyword test, not a parse: the engine deliberately contains a CAL
3829/// *writer*, never a parser (parsing is the substrate's job). Both spellings
3830/// are matched case-insensitively, and `DROP` is intentionally absent — the
3831/// loop may propose defining a query, never removing one.
3832pub(crate) fn is_definition_statement(line: &str) -> bool {
3833    let up = line.trim_start().to_ascii_uppercase();
3834    up.starts_with("DEFINE QUERY") || up.starts_with("DEFINE TEMPLATE")
3835}
3836
3837/// The refusal an advisory `Edit` earns. The engine has no executable edit
3838/// primitive; the change belongs in the host.
3839const ADVISORY_EDIT: &str = "This recommendation is advisory: the engine cannot execute this edit. \
3840     Make the change in the host. Dismiss it with REJECT … BECAUSE while it is still pending, or \
3841     approve it to acknowledge it and let it expire.";
3842
3843/// The refusal an advisory `Data` finding earns — every `Data` shape except
3844/// `outcome_review`'s revert, which carries `revert_of`.
3845/// Shared by [`Engine::preflight_apply`] and the apply gate so a fused
3846/// approve-and-apply caller is refused BEFORE the approval lands, not after.
3847const GATING_REQUIRED: &str = "code and adapter revisions apply only with a recorded gating run \
3848     (evalset hash + run id + stats) — use apply_gated";
3849
3850const ADVISORY_DATA: &str = "This finding is advisory: the engine cannot execute it. Act on its \
3851     guidance. Dismiss it with REJECT … BECAUSE while it is still pending, or approve it to \
3852     acknowledge it and let it expire.";
3853
3854/// Whether [`Engine::apply`] can execute this proposal at all.
3855///
3856/// One source of truth, shared by [`Engine::preflight_apply`] and
3857/// [`Engine::apply`] so the two can never disagree.
3858///
3859/// Deliberately **not** consulted on approve. An advisory finding is a `Flag`:
3860/// approving it means "yes, this is real", which is the whole workflow for the
3861/// LLM path and the telemetry analyzers. What must not happen is a caller being
3862/// walked into an approval and *then* refused — which is exactly what the fused
3863/// approve-and-apply path in the bindings did, leaving the recommendation in
3864/// `approved`, whose only exits are `applied` and `expired`. Preflight asks
3865/// first, so that path now refuses before it commits anything.
3866pub(crate) fn ensure_executable(action_kind: ActionKind, proposal: &Proposal) -> Result<()> {
3867    match proposal {
3868        Proposal::Cal { .. } => Ok(()),
3869        Proposal::Edit { .. } => Err(Error::InvalidProposal(ADVISORY_EDIT.into())),
3870        // Two executable Data shapes: `outcome_review`'s `revert_of` (names
3871        // an earlier applied recommendation to roll back), and a gated
3872        // revision (code or adapter), which executes by writing its
3873        // promotion grain — the gating-run requirement itself is checked at
3874        // apply, not here.
3875        Proposal::Data { data } => {
3876            if requires_gating(action_kind)
3877                || data.get("revert_of").and_then(Value::as_str).is_some()
3878            {
3879                Ok(())
3880            } else {
3881                Err(Error::InvalidProposal(ADVISORY_DATA.into()))
3882            }
3883        }
3884    }
3885}
3886
3887#[cfg(test)]
3888mod definition_body_tests {
3889    use super::safe_definition_body;
3890
3891    #[test]
3892    fn ordinary_bodies_pass() {
3893        assert!(safe_definition_body("RECALL facts WHERE relation = \"lesson\" LIMIT 20"));
3894        assert!(safe_definition_body("ASSEMBLE context FOR \"desk\" BUDGET 2000"));
3895    }
3896
3897    #[test]
3898    fn a_body_cannot_close_its_own_block_or_carry_destruction() {
3899        // The injection shape: close the DEFINE block, append a statement.
3900        // Newlines are already collapsed by `sanitize_line`, so the payload
3901        // arrives as ONE line — which is exactly what a line-leading-keyword
3902        // destructive scan cannot see.
3903        assert!(!safe_definition_body("RECALL facts } FORGET abc {"));
3904        assert!(!safe_definition_body("RECALL facts } PURGE OLDER THAN 1d {"));
3905        // …and the keyword alone is refused even without the braces.
3906        assert!(!safe_definition_body("RECALL facts FORGET abc"));
3907        assert!(!safe_definition_body("recall facts purge older than 1d"));
3908        assert!(!safe_definition_body("RECALL facts DROP QUERY \"x\""));
3909        // A nested DEFINE would redefine something the target does not name.
3910        assert!(!safe_definition_body("RECALL facts DEFINE QUERY \"other\""));
3911        // Substring matches are not keywords — this must still pass.
3912        assert!(safe_definition_body("RECALL facts WHERE subject = \"purged_at\""));
3913    }
3914}
3915
3916#[cfg(test)]
3917mod plan_edit_tests {
3918    use super::{plan_edit_allowed, plan_get, plan_set, plan_value_ok};
3919    use serde_json::json;
3920
3921    fn plan() -> serde_json::Value {
3922        json!({
3923            "nodes": ["fetch", "review", "post"],
3924            "edges": [
3925                {"src": "fetch", "dst": "review"},
3926                {"src": "review", "dst": "fetch", "cond": "confidence < 0.9", "max_cycles": 2}
3927            ],
3928            "bindings": {"fetch": "sha256:tool1"},
3929            "retries": {"fetch": 1}
3930        })
3931    }
3932
3933    #[test]
3934    fn the_allowlist_admits_thresholds_and_refuses_topology() {
3935        assert!(plan_edit_allowed("edges.1.cond"));
3936        assert!(plan_edit_allowed("edges.1.max_cycles"));
3937        assert!(plan_edit_allowed("retries.fetch"));
3938        // Topology is not expressible — the structural half of the guarantee
3939        // that a plan revision stays reviewable as scalar deltas.
3940        for path in [
3941            "nodes",
3942            "nodes.0",
3943            "edges.0.src",
3944            "edges.0.dst",
3945            "edges",
3946            "bindings.fetch",
3947            "edges.x.cond",
3948            "",
3949        ] {
3950            assert!(!plan_edit_allowed(path), "{path} must not be editable");
3951        }
3952    }
3953
3954    #[test]
3955    fn values_are_type_checked_against_the_field() {
3956        // A string in max_cycles would be DROPPED by the grain deserializer,
3957        // so an "applied" tightening would silently mean unlimited.
3958        assert!(!plan_value_ok("edges.1.max_cycles", &json!("2")));
3959        assert!(plan_value_ok("edges.1.max_cycles", &json!(2)));
3960        assert!(!plan_value_ok("edges.1.max_cycles", &json!(-1)));
3961        assert!(!plan_value_ok("retries.fetch", &json!(10_000)));
3962        assert!(plan_value_ok("retries.fetch", &json!(3)));
3963        assert!(plan_value_ok("edges.1.cond", &json!("confidence < 0.8")));
3964        assert!(!plan_value_ok("edges.1.cond", &json!("  ")));
3965        assert!(!plan_value_ok("edges.1.cond", &json!("a\nb")));
3966        assert!(!plan_value_ok("edges.0.src", &json!("other")));
3967    }
3968
3969    #[test]
3970    fn get_reads_through_arrays_and_objects_and_absence_is_null() {
3971        let p = plan();
3972        assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(2));
3973        assert_eq!(plan_get(&p, "retries.fetch"), json!(1));
3974        // An edit that ADDS a retry declares `from: null` — so absence has to
3975        // read as Null rather than as an error.
3976        assert_eq!(plan_get(&p, "retries.review"), json!(null));
3977        assert_eq!(plan_get(&p, "edges.9.cond"), json!(null));
3978        assert_eq!(plan_get(&p, "edges.0.cond"), json!(null));
3979    }
3980
3981    #[test]
3982    fn set_writes_scalars_and_adds_a_missing_retry_but_never_grows_an_array() {
3983        let mut p = plan();
3984        assert!(plan_set(&mut p, "edges.1.max_cycles", json!(5)));
3985        assert_eq!(plan_get(&p, "edges.1.max_cycles"), json!(5));
3986        assert!(plan_set(&mut p, "retries.review", json!(2)));
3987        assert_eq!(plan_get(&p, "retries.review"), json!(2));
3988        assert!(!plan_set(&mut p, "edges.7.cond", json!("x")));
3989        assert_eq!(p["edges"].as_array().unwrap().len(), 2, "no array growth");
3990    }
3991}
3992
3993#[cfg(test)]
3994mod definition_proposal_tests {
3995    use super::is_definition_statement;
3996
3997    #[test]
3998    fn definition_statements_are_recognized_in_both_spellings() {
3999        assert!(is_definition_statement(r#"DEFINE QUERY "triage" AS { RECALL facts }"#));
4000        assert!(is_definition_statement("  define template foo AS { x }"));
4001        assert!(is_definition_statement("DEFINE TEMPLATE bar AS { y }"));
4002        // Ordinary proposals are untouched.
4003        assert!(!is_definition_statement("ADD fact {}"));
4004        assert!(!is_definition_statement("SUPERSEDE abc WITH fact {}"));
4005        assert!(!is_definition_statement("FORGET abc"));
4006        // `DROP` is never proposable, so it is deliberately NOT a definition
4007        // statement here — a proposal containing one still fails validation
4008        // rather than being handed an inverse.
4009        assert!(!is_definition_statement(r#"DROP QUERY "triage""#));
4010    }
4011}