Skip to main content

areev_loop/
replay.rs

1//! Replay — score a loop configuration against the immutable past.
2//!
3//! The engine is a pure function of (file, policy, now): `run` never reads
4//! the clock, and the golden suite byte-pins queues because of it. So a
5//! candidate configuration has a measurable quality on the recorded past:
6//! which findings it would have produced at each historical pass, how many
7//! of those the humans went on to approve or reject, how many of the
8//! approved ones later regressed, and how much queue it would have made.
9//! Dream-RSI (arXiv 2609.14858) does this for exploration policies over
10//! discovery trees and always includes the incumbent in the candidate set;
11//! the same rule holds here — the report is a comparison, never a bare
12//! number. `docs/loop-proposal.md` §17 named this rung 1 of the ladder.
13//!
14//! Three disciplines, all structural:
15//!
16//! - **Prefix only.** Every step reads through [`PrefixView`], which hides
17//!   any grain created after the step's `now` — the paper's prefix rule, no
18//!   leakage from the future.
19//! - **Zero writes.** The view refuses every mutating method of
20//!   [`OmsSubstrate`], and the engine holds `&S`, not `&mut S`.
21//! - **Scope honesty.** An LLM is not a pure function of the evidence and
22//!   an external command is out of process, so neither is replayed; they
23//!   appear in the report as `not_replayed` with the reason. Telemetry
24//!   rollups are not time-indexed, so the telemetry-fed analyzers are not
25//!   replayed either.
26//!
27//! No auto-adoption: replay informs, and adopting the configuration remains
28//! the policy file or `POST /api/loop/config`, by a human.
29
30use crate::analyzer::OutcomeInput;
31use crate::config::{AnalyzerConfig, LoopPersisted};
32use crate::engine::{Engine, RunOptions, LOOP_NS};
33use crate::error::{Error, Result};
34use crate::model::{GrainRecord, Origin};
35use crate::policy::Policy;
36use crate::recommendation::{RecStatus, Recommendation};
37use crate::substrate::{Capabilities, GrainSpec, HeadGroup, OmsSubstrate, ReadOpts, SubstrateRead, TelemetryView};
38use serde::{Deserialize, Serialize};
39use serde_json::Value;
40use std::collections::{BTreeMap, BTreeSet};
41
42/// A candidate loop configuration to score against the past.
43#[derive(Debug, Clone, Default, Deserialize)]
44#[serde(deny_unknown_fields)]
45pub struct ReplayCandidate {
46    /// Per-analyzer configuration, keyed by full analyzer id — the same
47    /// shape `set_analyzer_config` stores (`enabled`, `params`,
48    /// `severity_floor`, `namespaces`). Each entry REPLACES the file's entry
49    /// for that analyzer; analyzers not named keep the file's config.
50    #[serde(default)]
51    pub config: BTreeMap<String, AnalyzerConfig>,
52    /// An optional host policy to replay under (severity floors, the deny
53    /// list, `near_duplicate`, …) instead of the engine's. Auto-apply
54    /// grants are irrelevant — a replay applies nothing.
55    #[serde(default)]
56    pub policy: Option<Policy>,
57}
58
59impl ReplayCandidate {
60    pub fn from_json(s: &str) -> Result<Self> {
61        serde_json::from_str(s).map_err(|e| Error::InvalidProposal(format!("replay candidate: {e}")))
62    }
63}
64
65/// How `now` steps through the window.
66#[derive(Debug, Clone, Copy, PartialEq, Eq)]
67pub enum ReplayStep {
68    /// One step per recorded pass — the moments the loop actually ran,
69    /// reconstructed from the audit trail (every stored finding is an
70    /// Observation stamped with the pass's `now`). A pass that stored
71    /// nothing left no trace and is not a step.
72    PerPass,
73    /// A fixed stride from the window's start, in ms.
74    Stride(i64),
75}
76
77impl ReplayStep {
78    /// `per-pass`, or a duration like `1d` / `12h` / `30m` / `86400000`.
79    pub fn parse(s: &str) -> Option<Self> {
80        let s = s.trim();
81        if s.eq_ignore_ascii_case("per-pass") || s.eq_ignore_ascii_case("per_pass") {
82            return Some(ReplayStep::PerPass);
83        }
84        let (num, unit) = match s.char_indices().find(|(_, c)| !c.is_ascii_digit()) {
85            Some((i, _)) => (&s[..i], &s[i..]),
86            None => (s, "ms"),
87        };
88        let n: i64 = num.parse().ok()?;
89        let mult = match unit {
90            "ms" => 1,
91            "s" => 1_000,
92            "m" => 60_000,
93            "h" => 3_600_000,
94            "d" => 86_400_000,
95            _ => return None,
96        };
97        (n > 0).then(|| ReplayStep::Stride(n * mult))
98    }
99
100    fn label(&self) -> String {
101        match self {
102            ReplayStep::PerPass => "per_pass".into(),
103            ReplayStep::Stride(ms) => format!("stride:{ms}ms"),
104        }
105    }
106}
107
108/// What to replay over.
109#[derive(Debug, Clone)]
110pub struct ReplayOptions {
111    /// The window's start (inclusive). `None` = from the first recorded pass
112    /// (per-pass) — a stride needs one.
113    pub since_ms: Option<i64>,
114    /// The window's end — the caller's `now`.
115    pub until_ms: i64,
116    pub step: ReplayStep,
117    /// The global namespace filter a live run would use (empty = all).
118    pub namespaces: Vec<String>,
119}
120
121impl ReplayOptions {
122    /// The one argument parser every surface shares: a `window` like `90d`
123    /// (back from `now`) or an explicit `since_ms`, and a `step` of
124    /// `per-pass` (default) or a duration. A window and a since together are
125    /// refused rather than silently ranked.
126    pub fn from_args(
127        window: Option<&str>,
128        since_ms: Option<i64>,
129        step: Option<&str>,
130        namespaces: Vec<String>,
131        now_ms: i64,
132    ) -> Result<Self> {
133        let bad = |what: String| Error::InvalidProposal(format!("replay: {what}"));
134        let since = match (window, since_ms) {
135            (Some(_), Some(_)) => return Err(bad("give a window or a since, not both".into())),
136            (Some(w), None) => match ReplayStep::parse(w) {
137                Some(ReplayStep::Stride(ms)) => Some(now_ms - ms),
138                _ => return Err(bad(format!("window {w:?} is not a duration like 90d, 12h or 30m"))),
139            },
140            (None, s) => s,
141        };
142        let step = match step {
143            None => ReplayStep::PerPass,
144            Some(s) => ReplayStep::parse(s)
145                .ok_or_else(|| bad(format!("step {s:?} is not per-pass or a duration like 1d")))?,
146        };
147        Ok(ReplayOptions { since_ms: since, until_ms: now_ms, step, namespaces })
148    }
149}
150
151/// The request every surface accepts (`areev loop replay --config FILE`
152/// reads the file into it and adds the flags; `POST /api/loop/replay` and
153/// the bindings take it whole): the candidate plus the window.
154#[derive(Debug, Clone, Default, Deserialize)]
155#[serde(deny_unknown_fields)]
156pub struct ReplayRequest {
157    #[serde(default)]
158    pub config: BTreeMap<String, AnalyzerConfig>,
159    #[serde(default)]
160    pub policy: Option<Policy>,
161    /// A duration back from now, e.g. `90d`.
162    #[serde(default)]
163    pub window: Option<String>,
164    #[serde(default)]
165    pub since_ms: Option<i64>,
166    /// `per-pass` (default) or a duration like `1d`.
167    #[serde(default)]
168    pub step: Option<String>,
169    #[serde(default)]
170    pub namespaces: Vec<String>,
171}
172
173impl ReplayRequest {
174    pub fn from_json(s: &str) -> Result<Self> {
175        serde_json::from_str(s).map_err(|e| Error::InvalidProposal(format!("replay request: {e}")))
176    }
177
178    pub fn resolve(self, now_ms: i64) -> Result<(ReplayCandidate, ReplayOptions)> {
179        let opts = ReplayOptions::from_args(
180            self.window.as_deref(),
181            self.since_ms,
182            self.step.as_deref(),
183            self.namespaces,
184            now_ms,
185        )?;
186        Ok((ReplayCandidate { config: self.config, policy: self.policy }, opts))
187    }
188}
189
190/// One would-be finding at one step.
191#[derive(Debug, Clone, Serialize, PartialEq)]
192pub struct ReplayFinding {
193    pub step_ms: i64,
194    pub analyzer: String,
195    pub dedup_key: String,
196    pub summary: String,
197    pub severity: String,
198    pub target_ref: String,
199    /// The content address a live pass would have stored it under, when the
200    /// substrate can compute one without writing.
201    #[serde(skip_serializing_if = "Option::is_none")]
202    pub address: Option<String>,
203    /// How the recorded history treated a finding with this dedup key:
204    /// `approved` (incl. applied / rolled back), `rejected`,
205    /// `never_reviewed` (stored, still pending or expired), or
206    /// `never_proposed` (the incumbent never produced it).
207    pub recorded: String,
208    /// The latest Verify-gate verdict on the recorded apply of this key, when
209    /// there was one: `held`, `held_costlier`, `regressed`, `drifted`.
210    #[serde(skip_serializing_if = "Option::is_none")]
211    pub outcome: Option<String>,
212}
213
214/// Counts for one analyzer, or the total.
215#[derive(Debug, Clone, Default, Serialize, PartialEq)]
216pub struct ReplayTally {
217    pub findings: u64,
218    pub approved: u64,
219    pub rejected: u64,
220    pub never_reviewed: u64,
221    pub never_proposed: u64,
222    pub regressed: u64,
223    pub drifted: u64,
224    pub held: u64,
225}
226
227impl ReplayTally {
228    fn add(&mut self, f: &ReplayFinding) {
229        self.findings += 1;
230        match f.recorded.as_str() {
231            "approved" => self.approved += 1,
232            "rejected" => self.rejected += 1,
233            "never_reviewed" => self.never_reviewed += 1,
234            _ => self.never_proposed += 1,
235        }
236        match f.outcome.as_deref() {
237            Some("regressed") => self.regressed += 1,
238            Some("drifted") => self.drifted += 1,
239            Some("held") | Some("held_costlier") => self.held += 1,
240            _ => {}
241        }
242    }
243}
244
245/// One arm of the comparison: the incumbent or the candidate.
246#[derive(Debug, Clone, Default, Serialize, PartialEq)]
247pub struct ReplayArm {
248    pub total: ReplayTally,
249    pub per_analyzer: BTreeMap<String, ReplayTally>,
250    /// Findings per step, in step order — the queue volume.
251    pub queue_per_step: Vec<u64>,
252    pub findings: Vec<ReplayFinding>,
253    /// Analyzers skipped at any step and why (disabled, denied, not replayed).
254    pub skipped: BTreeMap<String, String>,
255}
256
257/// Something the replay did not rehearse, and why.
258#[derive(Debug, Clone, Serialize, PartialEq)]
259pub struct NotReplayed {
260    pub what: String,
261    pub reason: String,
262}
263
264/// The report: a comparison, never a bare number.
265#[derive(Debug, Clone, Serialize)]
266pub struct ReplayReport {
267    pub since_ms: i64,
268    pub until_ms: i64,
269    pub step: String,
270    /// The `now` values replayed, in order.
271    pub steps: Vec<i64>,
272    pub incumbent: ReplayArm,
273    pub candidate: ReplayArm,
274    pub not_replayed: Vec<NotReplayed>,
275    /// Recorded review decisions and outcomes are matched to would-be
276    /// findings by dedup key — stated so the overlap columns are read for
277    /// what they are.
278    pub matching: &'static str,
279}
280
281/// A read-only, prefix-bounded view of a substrate. Every grain created
282/// after `until_ms` is invisible, and every write is refused — the two
283/// properties a rehearsal needs to be exact and harmless, both enforced by
284/// the type rather than promised.
285pub struct PrefixView<'a, S: OmsSubstrate> {
286    inner: &'a S,
287    until_ms: i64,
288}
289
290impl<'a, S: OmsSubstrate> PrefixView<'a, S> {
291    pub fn new(inner: &'a S, until_ms: i64) -> Self {
292        PrefixView { inner, until_ms }
293    }
294}
295
296fn read_only<T>(what: &str) -> Result<T> {
297    Err(Error::Substrate(format!("replay is read-only: {what} refused")))
298}
299
300impl<S: OmsSubstrate> SubstrateRead for PrefixView<'_, S> {
301    fn capabilities(&self) -> Capabilities {
302        self.inner.capabilities()
303    }
304    fn grains_of_type(&self, grain_type: &str, namespace: Option<&str>, opts: ReadOpts) -> Result<Vec<GrainRecord>> {
305        Ok(self
306            .inner
307            .grains_of_type(grain_type, namespace, opts)?
308            .into_iter()
309            .filter(|g| g.created_at_ms <= self.until_ms)
310            .collect())
311    }
312    fn grain(&self, hash: &str) -> Result<Option<GrainRecord>> {
313        Ok(self.inner.grain(hash)?.filter(|g| g.created_at_ms <= self.until_ms))
314    }
315    fn heads(&self, namespace: Option<&str>) -> Result<Vec<HeadGroup>> {
316        self.inner.heads(namespace)
317    }
318    fn telemetry(&self, namespace: Option<&str>) -> Result<Option<TelemetryView>> {
319        self.inner.telemetry(namespace)
320    }
321    fn validate_plan(&self, workflow: &Value) -> Result<()> {
322        self.inner.validate_plan(workflow)
323    }
324    fn tool_evalset(&self, tool: &str) -> Result<Option<String>> {
325        self.inner.tool_evalset(tool)
326    }
327    fn embed(&self, text: &str) -> Result<Option<Vec<f32>>> {
328        self.inner.embed(text)
329    }
330    fn address_of(&self, spec: &GrainSpec) -> Result<Option<String>> {
331        self.inner.address_of(spec)
332    }
333}
334
335impl<S: OmsSubstrate> OmsSubstrate for PrefixView<'_, S> {
336    fn put_grain(&mut self, _spec: &GrainSpec) -> Result<String> {
337        read_only("put_grain")
338    }
339    fn supersede(&mut self, _target_hash: &str, _spec: &GrainSpec, _justification: &str) -> Result<String> {
340        read_only("supersede")
341    }
342    fn retract(&mut self, _hash: &str, _reason: &str) -> Result<()> {
343        read_only("retract")
344    }
345    fn put_blob(&mut self, _bytes: &[u8]) -> Result<String> {
346        read_only("put_blob")
347    }
348    fn execute_cal(&mut self, _cal: &str) -> Result<Vec<Value>> {
349        read_only("execute_cal")
350    }
351    fn validate_cal(&self, cal: &str) -> Result<()> {
352        self.inner.validate_cal(cal)
353    }
354    fn definition_inverse(&self, statement: &str) -> Result<Option<String>> {
355        self.inner.definition_inverse(statement)
356    }
357    fn load_state(&self) -> Result<Value> {
358        self.inner.load_state()
359    }
360    fn store_state(&mut self, _state: &Value) -> Result<()> {
361        read_only("store_state")
362    }
363}
364
365/// One recorded transition on the audit trail.
366struct AuditEvent {
367    at_ms: i64,
368    rec_hash: String,
369    to: String,
370    from: Option<String>,
371}
372
373/// Recorded recommendations by hash, and the dedup-key index over them.
374type RecordedIndex = (BTreeMap<String, Recorded>, BTreeMap<String, Vec<String>>);
375
376/// A recorded recommendation, as the replay scores against it.
377struct Recorded {
378    dedup_key: String,
379    status: RecStatus,
380    outcome: Option<String>,
381    origin: Origin,
382    /// For a revert recommendation: the recommendation it retracts.
383    revert_of: Option<String>,
384}
385
386impl Engine {
387    /// Score `candidate` against the recorded past, beside the incumbent
388    /// (the file's config under this engine's policy). Writes nothing —
389    /// see the module docs for the three structural disciplines.
390    pub fn replay<S: OmsSubstrate>(
391        &self,
392        sub: &S,
393        candidate: &ReplayCandidate,
394        opts: &ReplayOptions,
395    ) -> Result<ReplayReport> {
396        let persisted = LoopPersisted::from_value(sub.load_state()?)?;
397        let (recorded, by_key) = self.recorded(sub, &persisted)?;
398        let events = audit_events(sub)?;
399
400        // The moments the loop ran, reconstructed from the trail: every
401        // stored finding's first transition is stamped with its pass's now.
402        let mut passes: BTreeSet<i64> = events
403            .iter()
404            .filter(|e| e.to == "pending" && e.from.is_none())
405            .map(|e| e.at_ms)
406            .collect();
407        if let Some(last) = persisted.state.last_run_ms {
408            passes.insert(last);
409        }
410        let since = match (opts.since_ms, passes.iter().next()) {
411            (Some(s), _) => s,
412            (None, Some(first)) => *first,
413            (None, None) => {
414                return Err(Error::InvalidProposal(
415                    "no recorded passes to replay through — give --since (or a window) and a --step".into(),
416                ))
417            }
418        };
419        let steps: Vec<i64> = match opts.step {
420            ReplayStep::PerPass => passes.into_iter().filter(|t| *t >= since && *t <= opts.until_ms).collect(),
421            ReplayStep::Stride(ms) => {
422                let mut v = Vec::new();
423                let mut t = since;
424                while t <= opts.until_ms {
425                    v.push(t);
426                    t += ms;
427                }
428                v
429            }
430        };
431        if steps.is_empty() {
432            return Err(Error::InvalidProposal(
433                "the window holds no step — widen it, or use a stride".into(),
434            ));
435        }
436
437        let mut not_replayed = Vec::new();
438        if self.has_llm() {
439            not_replayed.push(NotReplayed {
440                what: "origin=llm (the attached backend)".into(),
441                reason: "a model is not a pure function of the evidence".into(),
442            });
443        }
444        let llm_recorded = recorded.values().filter(|r| matches!(r.origin, Origin::Llm { .. })).count();
445        if llm_recorded > 0 {
446            not_replayed.push(NotReplayed {
447                what: format!("origin=llm ({llm_recorded} recorded finding(s))"),
448                reason: "a model is not a pure function of the evidence".into(),
449            });
450        }
451        let cmd_recorded = recorded.values().filter(|r| matches!(r.origin, Origin::Command { .. })).count();
452        for a in self.analyzers() {
453            let m = a.manifest();
454            if m.trust_class == crate::manifest::TrustClass::Command {
455                not_replayed.push(NotReplayed {
456                    what: format!("origin=command ({})", m.id),
457                    reason: "an external command is out of process".into(),
458                });
459            }
460        }
461        if cmd_recorded > 0 && !not_replayed.iter().any(|n| n.what.starts_with("origin=command")) {
462            not_replayed.push(NotReplayed {
463                what: format!("origin=command ({cmd_recorded} recorded finding(s))"),
464                reason: "an external command is out of process".into(),
465            });
466        }
467
468        // The incumbent: the file's config under this engine's policy. The
469        // candidate: its overlay on the file's config, under its own policy
470        // when it names one.
471        let mut candidate_config = persisted.config.clone();
472        for (id, cfg) in &candidate.config {
473            candidate_config.insert(id.clone(), cfg.clone());
474        }
475        let candidate_policy = candidate.policy.as_ref().unwrap_or(self.policy());
476        let incumbent = self.replay_arm(sub, &persisted, &persisted.config, self.policy(), &steps, opts, &events, &recorded, &by_key)?;
477        let candidate_arm = self.replay_arm(sub, &persisted, &candidate_config, candidate_policy, &steps, opts, &events, &recorded, &by_key)?;
478
479        Ok(ReplayReport {
480            since_ms: since,
481            until_ms: opts.until_ms,
482            step: opts.step.label(),
483            steps,
484            incumbent,
485            candidate: candidate_arm,
486            not_replayed,
487            matching: "recorded review decisions and Verify-gate outcomes are matched to would-be findings by dedup key",
488        })
489    }
490
491    /// The recorded recommendations, keyed by hash, plus a dedup-key index.
492    fn recorded<S: OmsSubstrate>(
493        &self,
494        sub: &S,
495        persisted: &LoopPersisted,
496    ) -> Result<RecordedIndex> {
497        let grains = sub.grains_of_type(
498            crate::model::grain_type::RECOMMENDATION,
499            Some(LOOP_NS),
500            ReadOpts { live_only: false, since_ms: None },
501        )?;
502        let mut recorded = BTreeMap::new();
503        let mut by_key: BTreeMap<String, Vec<String>> = BTreeMap::new();
504        for g in grains {
505            let Ok(rec) = Recommendation::from_fields(&g.hash, &g.fields) else { continue };
506            let status = persisted.status_index.get(&g.hash).copied().unwrap_or(RecStatus::Pending);
507            let outcome = persisted
508                .outcomes
509                .get(&g.hash)
510                .and_then(|v| v.iter().max_by_key(|o| o.measured_at_ms))
511                .map(|o| o.verdict.clone());
512            let revert_of = match &rec.proposal {
513                crate::recommendation::Proposal::Data { data } => {
514                    data.get("revert_of").and_then(Value::as_str).map(str::to_string)
515                }
516                _ => None,
517            };
518            by_key.entry(rec.dedup_key.clone()).or_default().push(g.hash.clone());
519            recorded.insert(
520                g.hash.clone(),
521                Recorded { dedup_key: rec.dedup_key, status, outcome, origin: rec.origin, revert_of },
522            );
523        }
524        Ok((recorded, by_key))
525    }
526
527    /// Walk the steps under one configuration, carrying the state the loop
528    /// had at each: the watermark advances per step, a recorded rejection
529    /// (or a measured revert) of a key puts it on the same doubling cooldown
530    /// a live pass would have, a rollback frees it, and the queue the
531    /// rehearsal itself produced is what dedups the next step.
532    #[allow(clippy::too_many_arguments)]
533    fn replay_arm<S: OmsSubstrate>(
534        &self,
535        sub: &S,
536        persisted: &LoopPersisted,
537        config: &BTreeMap<String, AnalyzerConfig>,
538        policy: &Policy,
539        steps: &[i64],
540        opts: &ReplayOptions,
541        events: &[AuditEvent],
542        recorded: &BTreeMap<String, Recorded>,
543        by_key: &BTreeMap<String, Vec<String>>,
544    ) -> Result<ReplayArm> {
545        let mut scratch = persisted.clone();
546        scratch.config = config.clone();
547        scratch.cooldowns.clear();
548        scratch.cooldown_strikes.clear();
549        let mut open: BTreeSet<String> = BTreeSet::new();
550        let mut arm = ReplayArm::default();
551        let run_opts = RunOptions {
552            namespaces: opts.namespaces.clone(),
553            ..RunOptions::default()
554        };
555        let no_outcomes: Vec<OutcomeInput> = Vec::new();
556        let mut prev: Option<i64> = None;
557        for &t in steps {
558            // Decisions recorded since the previous step, in order.
559            for e in events.iter().filter(|e| prev.is_none_or(|p| e.at_ms > p) && e.at_ms <= t) {
560                let Some(r) = recorded.get(&e.rec_hash) else { continue };
561                match e.to.as_str() {
562                    "rejected" => {
563                        open.remove(&r.dedup_key);
564                        crate::engine::strike_cooldown(&mut scratch, r.dedup_key.clone(), e.at_ms);
565                    }
566                    "rolled_back" => {
567                        open.remove(&r.dedup_key);
568                    }
569                    // Applying a revert is a verdict on the reverted finding:
570                    // the live path puts it on cooldown too.
571                    "applied" => {
572                        if let Some(target) = r.revert_of.as_ref().and_then(|h| recorded.get(h)) {
573                            open.remove(&target.dedup_key);
574                            crate::engine::strike_cooldown(&mut scratch, target.dedup_key.clone(), e.at_ms);
575                        }
576                    }
577                    _ => {}
578                }
579            }
580            let view = PrefixView::new(sub, t);
581            let pass = self.analysis_pass_inner(
582                &view,
583                &scratch,
584                policy,
585                &run_opts,
586                &BTreeMap::new(),
587                prev,
588                t,
589                &no_outcomes,
590                &open,
591                Some("replay"),
592            )?;
593            for sk in &pass.analyzers_skipped {
594                arm.skipped.entry(sk.id.clone()).or_insert_with(|| sk.reason.clone());
595            }
596            arm.queue_per_step.push(pass.survivors.len() as u64);
597            for rec in pass.survivors {
598                open.insert(rec.dedup_key.clone());
599                let (recorded_as, outcome) = match by_key.get(&rec.dedup_key) {
600                    None => ("never_proposed", None),
601                    Some(hashes) => {
602                        let rs: Vec<&Recorded> = hashes.iter().filter_map(|h| recorded.get(h)).collect();
603                        let approved = rs.iter().any(|r| {
604                            matches!(r.status, RecStatus::Approved | RecStatus::Applied | RecStatus::RolledBack)
605                        });
606                        let rejected = rs.iter().any(|r| r.status == RecStatus::Rejected);
607                        let outcome = rs.iter().filter_map(|r| r.outcome.clone()).next_back();
608                        if approved {
609                            ("approved", outcome)
610                        } else if rejected {
611                            ("rejected", None)
612                        } else {
613                            ("never_reviewed", None)
614                        }
615                    }
616                };
617                let address = rec
618                    .to_grain_spec(LOOP_NS)
619                    .ok()
620                    .and_then(|spec| sub.address_of(&spec).ok().flatten());
621                let f = ReplayFinding {
622                    step_ms: t,
623                    analyzer: rec.analyzer.clone(),
624                    dedup_key: rec.dedup_key.clone(),
625                    summary: rec.summary.render(),
626                    severity: rec.severity.as_str().to_string(),
627                    target_ref: rec.target_ref.clone(),
628                    address,
629                    recorded: recorded_as.into(),
630                    outcome,
631                };
632                arm.total.add(&f);
633                arm.per_analyzer.entry(rec.analyzer.clone()).or_default().add(&f);
634                arm.findings.push(f);
635            }
636            prev = Some(t);
637        }
638        Ok(arm)
639    }
640}
641
642/// Every audit transition on the trail, oldest first.
643fn audit_events<S: SubstrateRead>(sub: &S) -> Result<Vec<AuditEvent>> {
644    let obs = sub.grains_of_type(
645        crate::model::grain_type::OBSERVATION,
646        Some(LOOP_NS),
647        ReadOpts { live_only: false, since_ms: None },
648    )?;
649    let mut out: Vec<AuditEvent> = obs
650        .iter()
651        .filter(|g| g.str_field("observation_kind") == Some("loop_audit"))
652        .filter_map(|g| {
653            Some(AuditEvent {
654                at_ms: g.fields.get("at_ms").and_then(Value::as_i64).unwrap_or(g.created_at_ms),
655                rec_hash: g.str_field("rec_hash")?.to_string(),
656                to: g.str_field("to_status")?.to_string(),
657                from: g.str_field("from_status").map(str::to_string),
658            })
659        })
660        .collect();
661    out.sort_by(|a, b| a.at_ms.cmp(&b.at_ms).then(a.rec_hash.cmp(&b.rec_hash)));
662    Ok(out)
663}
664
665#[cfg(test)]
666mod tests {
667    use super::*;
668    use crate::engine::{Decision, ScopeSet};
669    use crate::recommendation::ObserverType;
670    use crate::testkit::TestSubstrate;
671    use serde_json::json;
672
673    const DAY: i64 = 86_400_000;
674
675    fn seeded() -> TestSubstrate {
676        let mut sub = TestSubstrate::new();
677        for i in 0..5 {
678            sub.add_tool_call_at("stripe_refund", true, "rate limited # retry later", 1_000 + i);
679        }
680        sub.add_fact_at("agent", "sam", "lives_in", "berlin", 2_000);
681        sub.add_fact_at("agent", "sam", "lives_in", "tokyo", 2_001);
682        sub
683    }
684
685    fn keys(findings: &[ReplayFinding]) -> BTreeSet<(String, String)> {
686        findings.iter().map(|f| (f.analyzer.clone(), f.summary.clone())).collect()
687    }
688
689    fn per_pass(since: Option<i64>, now: i64) -> ReplayOptions {
690        ReplayOptions { since_ms: since, until_ms: now, step: ReplayStep::PerPass, namespaces: vec![] }
691    }
692
693    /// Identity: the candidate equal to the current config, stepped through
694    /// the recorded passes, reproduces the queue the live passes stored —
695    /// same analyzers, same summaries, same count — and the incumbent row is
696    /// always present.
697    #[test]
698    fn replaying_the_incumbent_reproduces_the_recorded_queue() {
699        let mut sub = seeded();
700        let e = Engine::with_builtins();
701        e.run(&mut sub.inner, &RunOptions::default(), 10_000).unwrap();
702        let live: BTreeSet<(String, String)> = e
703            .recommendations(&sub.inner, None)
704            .unwrap()
705            .iter()
706            .map(|r| (r.analyzer.clone(), r.summary.render()))
707            .collect();
708        assert!(!live.is_empty());
709        let report = e.replay(&sub.inner, &ReplayCandidate::default(), &per_pass(None, 10_000)).unwrap();
710        assert_eq!(report.steps, vec![10_000]);
711        assert_eq!(keys(&report.incumbent.findings), live);
712        assert_eq!(keys(&report.candidate.findings), live, "an empty overlay IS the incumbent");
713        assert_eq!(report.incumbent.total.findings, live.len() as u64);
714        assert_eq!(report.incumbent.queue_per_step, vec![live.len() as u64]);
715        assert!(report.incumbent.findings.iter().all(|f| f.recorded == "never_reviewed"), "stored, still pending");
716        assert!(report.not_replayed.is_empty(), "no model, no command: nothing to disclaim");
717    }
718
719    /// A parameter change moves the queue as designed, and the decision
720    /// overlap moves with it.
721    #[test]
722    fn a_param_change_moves_the_queue_and_the_decision_overlap() {
723        let mut sub = seeded();
724        let e = Engine::with_builtins();
725        e.run(&mut sub.inner, &RunOptions::default(), 10_000).unwrap();
726        // Approve the tool-failure lesson, so the overlap has an approval.
727        let rec = e
728            .recommendations(&sub.inner, None)
729            .unwrap()
730            .into_iter()
731            .find(|r| r.analyzer.starts_with("loop.tool_failure"))
732            .unwrap();
733        e.review(&mut sub.inner, &rec.hash, Decision::Approve, "user:a", ObserverType::Human, &ScopeSet::all(), "ok", 10_500)
734            .unwrap();
735        // The candidate raises the cluster threshold past the seeded five.
736        let mut cand = ReplayCandidate::default();
737        cand.config.insert(
738            "loop.tool_failure/1".into(),
739            AnalyzerConfig { params: json!({"min_count": 50}).as_object().unwrap().clone(), ..Default::default() },
740        );
741        let report = e.replay(&sub.inner, &cand, &per_pass(None, 11_000)).unwrap();
742        assert_eq!(report.incumbent.total.approved, 1, "{:?}", report.incumbent.total);
743        assert!(report.incumbent.findings.iter().any(|f| f.analyzer.starts_with("loop.tool_failure")));
744        assert!(
745            !report.candidate.findings.iter().any(|f| f.analyzer.starts_with("loop.tool_failure")),
746            "the raised floor drops the cluster: {:?}",
747            report.candidate.findings
748        );
749        assert_eq!(report.candidate.total.approved, 0);
750        assert_eq!(report.candidate.total.findings + 1, report.incumbent.total.findings);
751    }
752
753    /// Zero writes: the grain count and the loop state are untouched, and the
754    /// view refuses every write by type.
755    #[test]
756    fn replay_writes_nothing() {
757        let mut sub = seeded();
758        let e = Engine::with_builtins();
759        e.run(&mut sub.inner, &RunOptions::default(), 10_000).unwrap();
760        let grains_before = sub.inner.grains_of_type("recommendation", None, ReadOpts { live_only: false, since_ms: None }).unwrap().len()
761            + sub.inner.grains_of_type("observation", None, ReadOpts { live_only: false, since_ms: None }).unwrap().len()
762            + sub.inner.grains_of_type("fact", None, ReadOpts { live_only: false, since_ms: None }).unwrap().len();
763        let state_before = sub.inner.load_state().unwrap();
764        e.replay(&sub.inner, &ReplayCandidate::default(), &per_pass(None, 10_000)).unwrap();
765        let grains_after = sub.inner.grains_of_type("recommendation", None, ReadOpts { live_only: false, since_ms: None }).unwrap().len()
766            + sub.inner.grains_of_type("observation", None, ReadOpts { live_only: false, since_ms: None }).unwrap().len()
767            + sub.inner.grains_of_type("fact", None, ReadOpts { live_only: false, since_ms: None }).unwrap().len();
768        assert_eq!(grains_before, grains_after);
769        assert_eq!(sub.inner.load_state().unwrap(), state_before);
770        let mut view = PrefixView::new(&sub.inner, 10_000);
771        assert!(view.put_grain(&GrainSpec::new("fact", "x")).is_err());
772        assert!(view.store_state(&json!({})).is_err());
773        assert!(view.execute_cal("ADD fact {}").is_err());
774    }
775
776    /// No future leakage: a grain dated after step t cannot contribute to a
777    /// finding at t, and does at a later step.
778    #[test]
779    fn a_grain_after_the_step_is_invisible_at_that_step() {
780        let mut sub = TestSubstrate::new();
781        // The contradiction needs both values; the second lands after t1.
782        sub.add_fact_at("agent", "sam", "lives_in", "berlin", 2_000);
783        let e = Engine::with_builtins();
784        e.run(&mut sub.inner, &RunOptions::default(), 10_000).unwrap(); // pass 1: nothing
785        sub.add_fact_at("agent", "sam", "lives_in", "tokyo", 15_000);
786        e.run(&mut sub.inner, &RunOptions::default(), 20_000).unwrap(); // pass 2: the contradiction
787        // Per pass: the first pass stored nothing and left no trace, so the
788        // only recorded step is the second — and it finds the pair.
789        let report = e.replay(&sub.inner, &ReplayCandidate::default(), &per_pass(None, 20_000)).unwrap();
790        assert_eq!(report.steps, vec![20_000]);
791        assert_eq!(report.incumbent.queue_per_step, vec![1], "{:?}", report.incumbent.findings);
792        assert_eq!(report.incumbent.findings[0].step_ms, 20_000);
793        // A stride that steps at 12_000 sees berlin alone: still nothing —
794        // tokyo (created 15_000) is in the future of that step.
795        let opts = ReplayOptions { since_ms: Some(12_000), until_ms: 20_000, step: ReplayStep::Stride(8_000), namespaces: vec![] };
796        let report = e.replay(&sub.inner, &ReplayCandidate::default(), &opts).unwrap();
797        assert_eq!(report.steps, vec![12_000, 20_000]);
798        assert_eq!(report.incumbent.queue_per_step, vec![0, 1]);
799    }
800
801    /// State fidelity: a finding rejected at t is on cooldown at t+1, so the
802    /// rehearsal does not count it again; the watermark advances per step.
803    #[test]
804    fn a_recorded_rejection_puts_the_key_on_cooldown_for_the_next_step() {
805        let mut sub = seeded();
806        let e = Engine::with_builtins();
807        e.run(&mut sub.inner, &RunOptions::default(), 10_000).unwrap();
808        let rec = e
809            .recommendations(&sub.inner, None)
810            .unwrap()
811            .into_iter()
812            .find(|r| r.analyzer.starts_with("loop.contradiction_sweep"))
813            .unwrap();
814        e.review(&mut sub.inner, &rec.hash, Decision::Reject, "user:a", ObserverType::Human, &ScopeSet::all(), "no", 10_500)
815            .unwrap();
816        // A second live pass a day later: the rejection cools the finding down.
817        e.run(&mut sub.inner, &RunOptions::default(), 10_000 + DAY).unwrap();
818        // Replay through both passes: the contradiction is found once, at
819        // the first step, and counted `rejected`.
820        let opts = ReplayOptions { since_ms: Some(10_000), until_ms: 10_000 + DAY, step: ReplayStep::Stride(DAY), namespaces: vec![] };
821        let report = e.replay(&sub.inner, &ReplayCandidate::default(), &opts).unwrap();
822        let contradictions: Vec<&ReplayFinding> = report
823            .incumbent
824            .findings
825            .iter()
826            .filter(|f| f.analyzer.starts_with("loop.contradiction_sweep"))
827            .collect();
828        assert_eq!(contradictions.len(), 1, "{contradictions:?}");
829        assert_eq!((contradictions[0].step_ms, contradictions[0].recorded.as_str()), (10_000, "rejected"));
830        // And every other finding is counted once, not once per step: the
831        // replayed queue dedups its own later steps.
832        assert_eq!(report.incumbent.queue_per_step[1], 0, "{:?}", report.incumbent.findings);
833    }
834
835    struct Cmd;
836    impl crate::analyzer::Analyzer for Cmd {
837        fn manifest(&self) -> &crate::manifest::AnalyzerManifest {
838            use std::sync::OnceLock;
839            static M: OnceLock<crate::manifest::AnalyzerManifest> = OnceLock::new();
840            M.get_or_init(|| crate::manifest::AnalyzerManifest {
841                id: "acme.pii/1".into(),
842                title: "PII".into(),
843                description: "external".into(),
844                tier: crate::manifest::Tier::T0,
845                cadence: crate::manifest::CadenceClass::Fast,
846                requires: vec![],
847                target_classes: vec![crate::manifest::TargetClass::Memory],
848                auto_apply: crate::manifest::AutoApplyClass::Never,
849                trust_class: crate::manifest::TrustClass::Command,
850                params: vec![],
851                default_on: true,
852            })
853        }
854        fn analyze(&self, _ctx: &crate::analyzer::AnalyzeCtx) -> Result<Vec<crate::recommendation::RecDraft>> {
855            panic!("an out-of-process analyzer must never run in a replay")
856        }
857    }
858    struct Llm;
859    impl crate::llm::LlmBackend for Llm {
860        fn model(&self) -> &str {
861            "mock"
862        }
863        fn complete(&self, _r: &str) -> Result<String> {
864            panic!("a model must never be called in a replay")
865        }
866    }
867
868    /// Scope honesty: an attached model and a registered external analyzer
869    /// are disclaimed, never run, and the deterministic rows are unaffected.
870    #[test]
871    fn a_model_and_an_external_analyzer_are_reported_not_replayed() {
872        let mut sub = seeded();
873        let plain = Engine::with_builtins();
874        plain.run(&mut sub.inner, &RunOptions::default(), 10_000).unwrap();
875        let baseline = plain.replay(&sub.inner, &ReplayCandidate::default(), &per_pass(None, 10_000)).unwrap();
876        let mut e = Engine::with_builtins().with_llm(Box::new(Llm));
877        e.register(Box::new(Cmd));
878        let report = e.replay(&sub.inner, &ReplayCandidate::default(), &per_pass(None, 10_000)).unwrap();
879        let whats: Vec<&str> = report.not_replayed.iter().map(|n| n.what.as_str()).collect();
880        assert!(whats.iter().any(|w| w.starts_with("origin=llm")), "{whats:?}");
881        assert!(whats.iter().any(|w| w.contains("acme.pii/1")), "{whats:?}");
882        assert_eq!(report.incumbent.skipped.get("acme.pii/1").map(String::as_str), Some("not replayed: replay"));
883        assert_eq!(keys(&report.incumbent.findings), keys(&baseline.incumbent.findings));
884    }
885
886    #[test]
887    fn steps_and_windows_parse() {
888        assert_eq!(ReplayStep::parse("per-pass"), Some(ReplayStep::PerPass));
889        assert_eq!(ReplayStep::parse("1d"), Some(ReplayStep::Stride(DAY)));
890        assert_eq!(ReplayStep::parse("12h"), Some(ReplayStep::Stride(12 * 3_600_000)));
891        assert_eq!(ReplayStep::parse("0d"), None);
892        assert_eq!(ReplayStep::parse("soon"), None);
893        let o = ReplayOptions::from_args(Some("90d"), None, None, vec![], 100 * DAY).unwrap();
894        assert_eq!((o.since_ms, o.step), (Some(10 * DAY), ReplayStep::PerPass));
895        assert!(ReplayOptions::from_args(Some("90d"), Some(1), None, vec![], 0).is_err());
896        assert!(ReplayOptions::from_args(None, None, Some("weekly"), vec![], 0).is_err());
897        let r = ReplayRequest::from_json(r#"{"config": {"loop.staleness/1": {"enabled": false}}, "window": "7d"}"#).unwrap();
898        let (c, o) = r.resolve(10 * DAY).unwrap();
899        assert_eq!(c.config.get("loop.staleness/1").and_then(|c| c.enabled), Some(false));
900        assert_eq!(o.since_ms, Some(3 * DAY));
901        assert!(ReplayRequest::from_json(r#"{"analyzers": {}}"#).is_err(), "unknown keys are refused");
902    }
903}