Skip to main content

areev_loop/
reference.rs

1//! An in-memory reference substrate: a grain map plus a deliberately naive CAL
2//! subset. Engine CI runs the full suite against it with zero Areev, so the
3//! portability claim stays testable, and it doubles as the conformance kit for
4//! third-party substrates (proposal §10).
5//!
6//! The CAL subset understands exactly the statements built-in analyzers emit
7//! (ADD / SUPERSEDE / FORGET / RETRACT) plus read verbs as no-ops. It is not a
8//! general CAL engine — a real substrate (Areev) provides that.
9
10use crate::error::{Error, Result};
11use crate::model::GrainRecord;
12use crate::substrate::{
13    Capabilities, GrainSpec, HeadGroup, OmsSubstrate, ReadOpts, SubstrateRead, TelemetryView,
14};
15use serde_json::{json, Map, Value};
16use std::collections::HashMap;
17
18#[derive(Default)]
19pub struct ReferenceSubstrate {
20    grains: Vec<GrainRecord>,
21    by_hash: HashMap<String, usize>,
22    caps: Capabilities,
23    mock_embedder: bool,
24    /// A canned `plan_replay` report, for tests of the loop's replay gate.
25    plan_replay: Option<Value>,
26    state: Value,
27    next_id: u64,
28    clock: i64,
29    /// Entity → competing head hashes (fork surfacing input).
30    heads_index: HashMap<String, Vec<String>>,
31    /// Injected recall-telemetry snapshot (turns on the `telemetry` capability).
32    telemetry: Option<TelemetryView>,
33    /// Saved definitions by key (`query:<name>` / `template:<name>`) → the
34    /// statement that would restore them. The reference implementation of the
35    /// host-metadata registry a real substrate keeps in the file.
36    definitions: HashMap<String, String>,
37}
38
39impl ReferenceSubstrate {
40    pub fn new() -> Self {
41        ReferenceSubstrate {
42            state: Value::Null,
43            // Declared because they are implemented below — the reference
44            // substrate is the conformance kit, so the seam it advertises is
45            // the seam an implementer must fill.
46            caps: Capabilities {
47                plans: true,
48                code: true,
49                ..Capabilities::default()
50            },
51            mock_embedder: false,
52            plan_replay: None,
53            ..Default::default()
54        }
55    }
56
57    /// Install a deterministic MOCK embedder: a hashed bag of words over a
58    /// small synonym table, so "record the vendor name" and "write down the
59    /// supplier's name" land close and "record the amount" does not. It
60    /// exists so the T1 near-duplicate leg can be exercised with no model;
61    /// it is not a semantic embedder and must never leave the test kit.
62    pub fn set_mock_embedder(&mut self) {
63        self.mock_embedder = true;
64        self.caps.embeddings = true;
65    }
66
67    /// Pretend the runtime rehearsed every plan revision and reported this
68    /// (`SubstrateRead::plan_replay`). The reference substrate has no
69    /// runtime, so a test of the replay gate hands it the report.
70    pub fn set_plan_replay(&mut self, report: Value) {
71        self.plan_replay = Some(report);
72    }
73
74    pub fn set_capabilities(&mut self, caps: Capabilities) {
75        self.caps = caps;
76    }
77
78    /// Register a fork (turns on the `forks` capability).
79    pub fn register_fork(&mut self, entity: &str, heads: &[&str]) {
80        self.caps.forks = true;
81        self.heads_index.insert(
82            entity.to_string(),
83            heads.iter().map(|s| s.to_string()).collect(),
84        );
85    }
86
87    /// Inject a telemetry snapshot (turns on the `telemetry` capability).
88    pub fn set_telemetry(&mut self, view: TelemetryView) {
89        self.caps.telemetry = true;
90        self.telemetry = Some(view);
91    }
92
93    /// Insert a fully-formed grain record; returns its assigned hash.
94    pub fn insert(&mut self, mut record: GrainRecord) -> String {
95        let hash = if record.hash.is_empty() {
96            self.mint_hash()
97        } else {
98            record.hash.clone()
99        };
100        record.hash = hash.clone();
101        let idx = self.grains.len();
102        self.by_hash.insert(hash.clone(), idx);
103        self.grains.push(record);
104        hash
105    }
106
107    fn mint_hash(&mut self) -> String {
108        let h = format!("ref-{:08}", self.next_id);
109        self.next_id += 1;
110        h
111    }
112
113    fn tick(&mut self) -> i64 {
114        self.clock += 1;
115        self.clock
116    }
117
118    fn record_from_spec(&mut self, spec: &GrainSpec) -> GrainRecord {
119        let created = self.tick();
120        let namespace = spec
121            .fields
122            .get("namespace")
123            .and_then(Value::as_str)
124            .unwrap_or(&spec.namespace)
125            .to_string();
126        let valid_to_ms = spec.fields.get("valid_to_ms").and_then(Value::as_i64);
127        GrainRecord {
128            hash: String::new(),
129            grain_type: spec.grain_type.clone(),
130            namespace,
131            created_at_ms: created,
132            valid_to_ms,
133            superseded_by: None,
134            fields: spec.fields.clone(),
135        }
136    }
137}
138
139impl SubstrateRead for ReferenceSubstrate {
140    fn capabilities(&self) -> Capabilities {
141        self.caps
142    }
143
144    fn embed(&self, text: &str) -> Result<Option<Vec<f32>>> {
145        if !self.mock_embedder {
146            return Ok(None);
147        }
148        Ok(Some(mock_embed(text)))
149    }
150
151    fn plan_replay(&self, _incumbent: &str, _candidate: &Value) -> Result<Option<Value>> {
152        Ok(self.plan_replay.clone())
153    }
154
155    /// A MINIMAL plan check: every node named once, every edge between known
156    /// nodes, at least one node. Deliberately weaker than the Areev adapter's,
157    /// which hands the body to the runtime's own `PlanGraph::build` — this
158    /// crate carries no Areev dependency, and the seam's contract is "refuse a
159    /// plan you cannot vouch for", not "reimplement the scheduler".
160    fn validate_plan(&self, workflow: &Value) -> Result<()> {
161        let bad = |m: &str| Err(Error::Substrate(m.to_string()));
162        let Some(fields) = workflow.as_object() else {
163            return bad("a workflow body must be a JSON object");
164        };
165        let Some(Value::Array(nodes)) = fields.get("nodes") else {
166            return bad("a workflow needs a 'nodes' array");
167        };
168        let mut seen = std::collections::BTreeSet::new();
169        for n in nodes {
170            match n.as_str() {
171                Some(id) if seen.insert(id) => {}
172                Some(id) => return bad(&format!("duplicate node {id:?}")),
173                None => return bad("every workflow node must be a string"),
174            }
175        }
176        if seen.is_empty() {
177            return bad("a workflow needs at least one node");
178        }
179        if let Some(v) = fields.get("edges") {
180            let Some(edges) = v.as_array() else {
181                return bad("workflow 'edges' must be an array");
182            };
183            for e in edges {
184                for end in ["src", "dst"] {
185                    match e.get(end).and_then(Value::as_str) {
186                        Some(id) if seen.contains(id) => {}
187                        Some(id) => return bad(&format!("edge {end} {id:?} is not a node")),
188                        None => return bad(&format!("every edge needs a '{end}'")),
189                    }
190                }
191            }
192        }
193        Ok(())
194    }
195
196    /// Rule E1's pin, from the tool's own live definition grain.
197    fn tool_evalset(&self, tool: &str) -> Result<Option<String>> {
198        Ok(self
199            .grains
200            .iter()
201            .filter(|g| {
202                g.grain_type == crate::model::grain_type::TOOL
203                    && g.is_live()
204                    && g.str_field("kind") == Some("definition")
205                    && g.tool_name() == Some(tool)
206            })
207            .find_map(|g| {
208                g.str_field("evalset_hash")
209                    .filter(|h| !h.trim().is_empty())
210                    .map(str::to_string)
211            }))
212    }
213
214    fn grains_of_type(
215        &self,
216        grain_type: &str,
217        namespace: Option<&str>,
218        opts: ReadOpts,
219    ) -> Result<Vec<GrainRecord>> {
220        Ok(self
221            .grains
222            .iter()
223            .filter(|g| g.grain_type == grain_type)
224            .filter(|g| namespace.is_none_or(|ns| g.namespace == ns))
225            // An ALL-NAMESPACE scan hides governance namespaces, matching what
226            // a real substrate must do: those hold a file's grants and audit
227            // records, and an analyzer sweeping them as ordinary memory can
228            // propose tombstoning the grants. An EXPLICITLY named namespace is
229            // the caller saying they mean it, and is served.
230            //
231            // Modelled here rather than left to the adapters because a fake
232            // more permissive than production hides exactly the bug it should
233            // catch: the fold-summary carve-out passed against this substrate
234            // while the real one returned an empty evidence bundle, and the
235            // gap only showed up against a live model.
236            .filter(|g| namespace.is_some() || !g.namespace.starts_with("agent:"))
237            .filter(|g| !opts.live_only || g.is_live())
238            .filter(|g| opts.since_ms.is_none_or(|s| g.created_at_ms >= s))
239            .cloned()
240            .collect())
241    }
242
243    fn grain(&self, hash: &str) -> Result<Option<GrainRecord>> {
244        Ok(self.by_hash.get(hash).map(|&i| self.grains[i].clone()))
245    }
246
247    fn heads(&self, _namespace: Option<&str>) -> Result<Vec<HeadGroup>> {
248        if !self.caps.forks {
249            return Err(Error::CapabilityMissing("forks".into()));
250        }
251        let mut groups: Vec<HeadGroup> = self
252            .heads_index
253            .iter()
254            .map(|(entity, heads)| HeadGroup {
255                entity: entity.clone(),
256                heads: heads.clone(),
257            })
258            .collect();
259        groups.sort_by(|a, b| a.entity.cmp(&b.entity));
260        Ok(groups)
261    }
262
263    fn telemetry(&self, _namespace: Option<&str>) -> Result<Option<TelemetryView>> {
264        Ok(self.telemetry.clone())
265    }
266}
267
268impl OmsSubstrate for ReferenceSubstrate {
269    fn put_grain(&mut self, spec: &GrainSpec) -> Result<String> {
270        let record = self.record_from_spec(spec);
271        Ok(self.insert(record))
272    }
273
274    fn supersede(
275        &mut self,
276        target_hash: &str,
277        spec: &GrainSpec,
278        _justification: &str,
279    ) -> Result<String> {
280        let record = self.record_from_spec(spec);
281        let new_hash = self.insert(record);
282        let idx = *self
283            .by_hash
284            .get(target_hash)
285            .ok_or_else(|| Error::Substrate(format!("supersede target {target_hash} not found")))?;
286        self.grains[idx].superseded_by = Some(new_hash.clone());
287        Ok(new_hash)
288    }
289
290    fn retract(&mut self, hash: &str, reason: &str) -> Result<()> {
291        let idx = *self
292            .by_hash
293            .get(hash)
294            .ok_or_else(|| Error::Substrate(format!("retract target {hash} not found")))?;
295        self.grains[idx].superseded_by = Some("retracted".to_string());
296        self.grains[idx]
297            .fields
298            .insert("verification_status".into(), json!("retracted"));
299        self.grains[idx]
300            .fields
301            .insert("retract_reason".into(), json!(reason));
302        // Retracting a grain that superseded others restores them as heads —
303        // the Areev store's semantics (a rolled-back consolidation puts every
304        // member back), mirrored here so engine tests see the same world.
305        for g in self.grains.iter_mut() {
306            if g.superseded_by.as_deref() == Some(hash) {
307                g.superseded_by = None;
308            }
309        }
310        Ok(())
311    }
312
313    fn execute_cal(&mut self, cal: &str) -> Result<Vec<Value>> {
314        let mut rows = Vec::new();
315        for line in cal.lines() {
316            let line = line.trim();
317            if line.is_empty() {
318                continue;
319            }
320            let (keyword, rest) = split_keyword(line);
321            match keyword.to_ascii_uppercase().as_str() {
322                "FORGET" => {
323                    let hash = rest.trim();
324                    if let Some(&idx) = self.by_hash.get(hash) {
325                        self.grains[idx].superseded_by = Some("forgotten".to_string());
326                    }
327                }
328                "RETRACT" => {
329                    let hash = rest.trim();
330                    self.retract(hash, "cal retract")?;
331                }
332                "ADD" => {
333                    let (grain_type, fields) = parse_type_and_json(rest)?;
334                    let spec = GrainSpec {
335                        grain_type,
336                        namespace: String::new(),
337                        fields,
338                    };
339                    let h = self.put_grain(&spec)?;
340                    rows.push(json!({ "hash": h }));
341                }
342                "SUPERSEDE" => {
343                    // SUPERSEDE <hash> WITH <type> {json}
344                    let (target, after_with) = rest.split_once(" WITH ").ok_or_else(|| {
345                        Error::CalUnsupported(format!("malformed SUPERSEDE: {line}"))
346                    })?;
347                    let (grain_type, fields) = parse_type_and_json(after_with)?;
348                    let spec = GrainSpec {
349                        grain_type,
350                        namespace: String::new(),
351                        fields,
352                    };
353                    let h = self.supersede(target.trim(), &spec, "cal supersede")?;
354                    rows.push(json!({ "hash": h }));
355                }
356                "DEFINE" => {
357                    let (key, _) = definition_key(line).ok_or_else(|| {
358                        Error::CalUnsupported(format!("malformed DEFINE: {line}"))
359                    })?;
360                    self.definitions.insert(key, line.to_string());
361                }
362                "DROP" => {
363                    if let Some((key, _)) = definition_key(line) {
364                        self.definitions.remove(&key);
365                    }
366                }
367                // Read verbs: no-ops in the reference substrate (no metric value).
368                "RECALL" | "ASSEMBLE" | "EXPLAIN" | "HISTORY" => {}
369                other => {
370                    return Err(Error::CalUnsupported(format!(
371                        "unknown statement {other:?}"
372                    )));
373                }
374            }
375        }
376        Ok(rows)
377    }
378
379    fn validate_cal(&self, cal: &str) -> Result<()> {
380        for line in cal.lines() {
381            let line = line.trim();
382            if line.is_empty() {
383                continue;
384            }
385            let (keyword, _) = split_keyword(line);
386            match keyword.to_ascii_uppercase().as_str() {
387                "ADD" | "SUPERSEDE" | "FORGET" | "RETRACT" | "RECALL" | "ASSEMBLE" | "EXPLAIN"
388                | "HISTORY" => {}
389                "DEFINE" | "DROP" => {
390                    if definition_key(line).is_none() {
391                        return Err(Error::CalUnsupported(format!(
392                            "malformed definition statement: {line}"
393                        )));
394                    }
395                }
396                other => {
397                    return Err(Error::CalUnsupported(format!(
398                        "unknown statement {other:?}"
399                    )))
400                }
401            }
402        }
403        Ok(())
404    }
405
406    /// The statement restoring the CURRENT definition, or a `DROP` when there
407    /// is none — so an apply can always record an inverse and a ROLLBACK
408    /// always undoes something.
409    fn definition_inverse(&self, statement: &str) -> Result<Option<String>> {
410        let Some((key, drop_stmt)) = definition_key(statement) else {
411            return Ok(None);
412        };
413        Ok(Some(
414            self.definitions.get(&key).cloned().unwrap_or(drop_stmt),
415        ))
416    }
417
418    fn load_state(&self) -> Result<Value> {
419        Ok(self.state.clone())
420    }
421
422    fn store_state(&mut self, state: &Value) -> Result<()> {
423        self.state = state.clone();
424        Ok(())
425    }
426}
427
428/// Parse a `DEFINE`/`DROP` statement into its registry key and the `DROP`
429/// that would remove it. `None` for anything else — this is a reader of two
430/// statement shapes, not a CAL parser.
431fn definition_key(line: &str) -> Option<(String, String)> {
432    let rest = {
433        let (kw, rest) = split_keyword(line.trim());
434        match kw.to_ascii_uppercase().as_str() {
435            "DEFINE" | "DROP" => rest,
436            _ => return None,
437        }
438    };
439    let (kind, rest) = split_keyword(rest);
440    let kind = kind.to_ascii_uppercase();
441    let (name, quoted) = match kind.as_str() {
442        "QUERY" => {
443            let rest = rest.trim().strip_prefix('"')?;
444            (rest.split('"').next()?, true)
445        }
446        "TEMPLATE" => (split_keyword(rest.trim()).0, false),
447        _ => return None,
448    };
449    if name.is_empty() {
450        return None;
451    }
452    let lower = kind.to_ascii_lowercase();
453    let drop_stmt = if quoted {
454        format!("DROP {kind} \"{name}\"")
455    } else {
456        format!("DROP {kind} {name}")
457    };
458    Some((format!("{lower}:{name}"), drop_stmt))
459}
460
461fn split_keyword(line: &str) -> (&str, &str) {
462    match line.split_once(char::is_whitespace) {
463        Some((k, rest)) => (k, rest.trim_start()),
464        None => (line, ""),
465    }
466}
467
468/// Parse `<type> {json}` → (type, fields).
469fn parse_type_and_json(s: &str) -> Result<(String, Map<String, Value>)> {
470    let brace = s
471        .find('{')
472        .ok_or_else(|| Error::CalUnsupported(format!("missing JSON object in {s:?}")))?;
473    let grain_type = s[..brace].trim().to_string();
474    if grain_type.is_empty() {
475        return Err(Error::CalUnsupported(format!(
476            "missing grain type in {s:?}"
477        )));
478    }
479    let value: Value = serde_json::from_str(s[brace..].trim())
480        .map_err(|e| Error::CalUnsupported(format!("bad JSON in {s:?}: {e}")))?;
481    let obj = value
482        .as_object()
483        .ok_or_else(|| Error::CalUnsupported(format!("JSON not an object in {s:?}")))?
484        .clone();
485    Ok((grain_type, obj))
486}
487
488/// The test kit's embedder: canonicalize tokens through a small synonym table,
489/// then hash each into one of 64 buckets. Deterministic, model-free, and only
490/// as "semantic" as the table — which is the point: it lets a test say "these
491/// two lines mean the same" without a network call.
492fn mock_embed(text: &str) -> Vec<f32> {
493    const SYNONYMS: &[(&str, &str)] = &[
494        ("write", "record"), ("note", "record"), ("log", "record"), ("capture", "record"),
495        ("down", ""), ("supplier", "vendor"), ("seller", "vendor"), ("merchant", "vendor"),
496        ("always", ""), ("every", "each"), ("all", "each"), ("a", ""), ("an", ""), ("the", ""),
497        ("on", ""), ("of", ""), ("s", ""), ("for", ""), ("to", ""), ("and", ""), ("with", ""),
498        ("before", "prior"), ("ahead", "prior"), ("answering", "answer"), ("answers", "answer"),
499        ("confirm", "check"), ("verify", "check"), ("current", "present"), ("latest", "present"),
500        ("city", "location"), ("town", "location"), ("place", "location"),
501        ("invoice", "bill"), ("receipt", "bill"), ("number", "id"), ("identifier", "id"),
502        ("exactly", "verbatim"), ("printed", "shown"),
503    ];
504    let mut v = vec![0f32; 64];
505    for raw in text.to_lowercase().split(|c: char| !c.is_alphanumeric()) {
506        if raw.is_empty() {
507            continue;
508        }
509        let tok = SYNONYMS.iter().find(|(from, _)| *from == raw).map(|(_, to)| *to).unwrap_or(raw);
510        if tok.is_empty() {
511            continue;
512        }
513        // FNV-1a, so the bucket is a pure function of the token.
514        let mut h: u64 = 0xcbf29ce484222325;
515        for b in tok.bytes() {
516            h ^= b as u64;
517            h = h.wrapping_mul(0x100000001b3);
518        }
519        v[(h % 64) as usize] += 1.0;
520    }
521    v
522}
523
524#[cfg(test)]
525mod tests {
526    use super::*;
527
528    #[test]
529    fn forget_makes_grain_not_live() {
530        let mut sub = ReferenceSubstrate::new();
531        let h = sub
532            .put_grain(&GrainSpec::new("fact", "ns").with_field("subject", "x"))
533            .unwrap();
534        assert_eq!(
535            sub.grains_of_type("fact", None, ReadOpts::default())
536                .unwrap()
537                .len(),
538            1
539        );
540        sub.execute_cal(&format!("FORGET {h}")).unwrap();
541        assert!(sub
542            .grains_of_type("fact", None, ReadOpts::default())
543            .unwrap()
544            .is_empty());
545    }
546
547    #[test]
548    fn add_returns_hash_and_stores() {
549        let mut sub = ReferenceSubstrate::new();
550        let rows = sub
551            .execute_cal(r#"ADD fact {"subject":"acme","relation":"tier","object":"ent"}"#)
552            .unwrap();
553        assert_eq!(rows.len(), 1);
554        assert!(rows[0].get("hash").is_some());
555        assert_eq!(
556            sub.grains_of_type("fact", None, ReadOpts::default())
557                .unwrap()
558                .len(),
559            1
560        );
561    }
562
563    #[test]
564    fn validate_rejects_unknown_statement() {
565        let sub = ReferenceSubstrate::new();
566        assert!(sub.validate_cal("DROP TABLE").is_err());
567        assert!(sub.validate_cal("ADD fact {}").is_ok());
568    }
569
570    #[test]
571    fn state_round_trips() {
572        let mut sub = ReferenceSubstrate::new();
573        assert!(sub.load_state().unwrap().is_null());
574        sub.store_state(&json!({"k": 1})).unwrap();
575        assert_eq!(sub.load_state().unwrap(), json!({"k": 1}));
576    }
577}