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    state: Value,
24    next_id: u64,
25    clock: i64,
26    /// Entity → competing head hashes (fork surfacing input).
27    heads_index: HashMap<String, Vec<String>>,
28    /// Injected recall-telemetry snapshot (turns on the `telemetry` capability).
29    telemetry: Option<TelemetryView>,
30    /// Saved definitions by key (`query:<name>` / `template:<name>`) → the
31    /// statement that would restore them. The reference implementation of the
32    /// host-metadata registry a real substrate keeps in the file.
33    definitions: HashMap<String, String>,
34}
35
36impl ReferenceSubstrate {
37    pub fn new() -> Self {
38        ReferenceSubstrate {
39            state: Value::Null,
40            // Declared because they are implemented below — the reference
41            // substrate is the conformance kit, so the seam it advertises is
42            // the seam an implementer must fill.
43            caps: Capabilities {
44                plans: true,
45                code: true,
46                ..Capabilities::default()
47            },
48            ..Default::default()
49        }
50    }
51
52    pub fn set_capabilities(&mut self, caps: Capabilities) {
53        self.caps = caps;
54    }
55
56    /// Register a fork (turns on the `forks` capability).
57    pub fn register_fork(&mut self, entity: &str, heads: &[&str]) {
58        self.caps.forks = true;
59        self.heads_index.insert(
60            entity.to_string(),
61            heads.iter().map(|s| s.to_string()).collect(),
62        );
63    }
64
65    /// Inject a telemetry snapshot (turns on the `telemetry` capability).
66    pub fn set_telemetry(&mut self, view: TelemetryView) {
67        self.caps.telemetry = true;
68        self.telemetry = Some(view);
69    }
70
71    /// Insert a fully-formed grain record; returns its assigned hash.
72    pub fn insert(&mut self, mut record: GrainRecord) -> String {
73        let hash = if record.hash.is_empty() {
74            self.mint_hash()
75        } else {
76            record.hash.clone()
77        };
78        record.hash = hash.clone();
79        let idx = self.grains.len();
80        self.by_hash.insert(hash.clone(), idx);
81        self.grains.push(record);
82        hash
83    }
84
85    fn mint_hash(&mut self) -> String {
86        let h = format!("ref-{:08}", self.next_id);
87        self.next_id += 1;
88        h
89    }
90
91    fn tick(&mut self) -> i64 {
92        self.clock += 1;
93        self.clock
94    }
95
96    fn record_from_spec(&mut self, spec: &GrainSpec) -> GrainRecord {
97        let created = self.tick();
98        let namespace = spec
99            .fields
100            .get("namespace")
101            .and_then(Value::as_str)
102            .unwrap_or(&spec.namespace)
103            .to_string();
104        let valid_to_ms = spec.fields.get("valid_to_ms").and_then(Value::as_i64);
105        GrainRecord {
106            hash: String::new(),
107            grain_type: spec.grain_type.clone(),
108            namespace,
109            created_at_ms: created,
110            valid_to_ms,
111            superseded_by: None,
112            fields: spec.fields.clone(),
113        }
114    }
115}
116
117impl SubstrateRead for ReferenceSubstrate {
118    fn capabilities(&self) -> Capabilities {
119        self.caps
120    }
121
122    /// A MINIMAL plan check: every node named once, every edge between known
123    /// nodes, at least one node. Deliberately weaker than the Areev adapter's,
124    /// which hands the body to the runtime's own `PlanGraph::build` — this
125    /// crate carries no Areev dependency, and the seam's contract is "refuse a
126    /// plan you cannot vouch for", not "reimplement the scheduler".
127    fn validate_plan(&self, workflow: &Value) -> Result<()> {
128        let bad = |m: &str| Err(Error::Substrate(m.to_string()));
129        let Some(fields) = workflow.as_object() else {
130            return bad("a workflow body must be a JSON object");
131        };
132        let Some(Value::Array(nodes)) = fields.get("nodes") else {
133            return bad("a workflow needs a 'nodes' array");
134        };
135        let mut seen = std::collections::BTreeSet::new();
136        for n in nodes {
137            match n.as_str() {
138                Some(id) if seen.insert(id) => {}
139                Some(id) => return bad(&format!("duplicate node {id:?}")),
140                None => return bad("every workflow node must be a string"),
141            }
142        }
143        if seen.is_empty() {
144            return bad("a workflow needs at least one node");
145        }
146        if let Some(v) = fields.get("edges") {
147            let Some(edges) = v.as_array() else {
148                return bad("workflow 'edges' must be an array");
149            };
150            for e in edges {
151                for end in ["src", "dst"] {
152                    match e.get(end).and_then(Value::as_str) {
153                        Some(id) if seen.contains(id) => {}
154                        Some(id) => return bad(&format!("edge {end} {id:?} is not a node")),
155                        None => return bad(&format!("every edge needs a '{end}'")),
156                    }
157                }
158            }
159        }
160        Ok(())
161    }
162
163    /// Rule E1's pin, from the tool's own live definition grain.
164    fn tool_evalset(&self, tool: &str) -> Result<Option<String>> {
165        Ok(self
166            .grains
167            .iter()
168            .filter(|g| {
169                g.grain_type == crate::model::grain_type::TOOL
170                    && g.is_live()
171                    && g.str_field("kind") == Some("definition")
172                    && g.tool_name() == Some(tool)
173            })
174            .find_map(|g| {
175                g.str_field("evalset_hash")
176                    .filter(|h| !h.trim().is_empty())
177                    .map(str::to_string)
178            }))
179    }
180
181    fn grains_of_type(
182        &self,
183        grain_type: &str,
184        namespace: Option<&str>,
185        opts: ReadOpts,
186    ) -> Result<Vec<GrainRecord>> {
187        Ok(self
188            .grains
189            .iter()
190            .filter(|g| g.grain_type == grain_type)
191            .filter(|g| namespace.is_none_or(|ns| g.namespace == ns))
192            // An ALL-NAMESPACE scan hides governance namespaces, matching what
193            // a real substrate must do: those hold a file's grants and audit
194            // records, and an analyzer sweeping them as ordinary memory can
195            // propose tombstoning the grants. An EXPLICITLY named namespace is
196            // the caller saying they mean it, and is served.
197            //
198            // Modelled here rather than left to the adapters because a fake
199            // more permissive than production hides exactly the bug it should
200            // catch: the fold-summary carve-out passed against this substrate
201            // while the real one returned an empty evidence bundle, and the
202            // gap only showed up against a live model.
203            .filter(|g| namespace.is_some() || !g.namespace.starts_with("agent:"))
204            .filter(|g| !opts.live_only || g.is_live())
205            .filter(|g| opts.since_ms.is_none_or(|s| g.created_at_ms >= s))
206            .cloned()
207            .collect())
208    }
209
210    fn grain(&self, hash: &str) -> Result<Option<GrainRecord>> {
211        Ok(self.by_hash.get(hash).map(|&i| self.grains[i].clone()))
212    }
213
214    fn heads(&self, _namespace: Option<&str>) -> Result<Vec<HeadGroup>> {
215        if !self.caps.forks {
216            return Err(Error::CapabilityMissing("forks".into()));
217        }
218        let mut groups: Vec<HeadGroup> = self
219            .heads_index
220            .iter()
221            .map(|(entity, heads)| HeadGroup {
222                entity: entity.clone(),
223                heads: heads.clone(),
224            })
225            .collect();
226        groups.sort_by(|a, b| a.entity.cmp(&b.entity));
227        Ok(groups)
228    }
229
230    fn telemetry(&self, _namespace: Option<&str>) -> Result<Option<TelemetryView>> {
231        Ok(self.telemetry.clone())
232    }
233}
234
235impl OmsSubstrate for ReferenceSubstrate {
236    fn put_grain(&mut self, spec: &GrainSpec) -> Result<String> {
237        let record = self.record_from_spec(spec);
238        Ok(self.insert(record))
239    }
240
241    fn supersede(
242        &mut self,
243        target_hash: &str,
244        spec: &GrainSpec,
245        _justification: &str,
246    ) -> Result<String> {
247        let record = self.record_from_spec(spec);
248        let new_hash = self.insert(record);
249        let idx = *self
250            .by_hash
251            .get(target_hash)
252            .ok_or_else(|| Error::Substrate(format!("supersede target {target_hash} not found")))?;
253        self.grains[idx].superseded_by = Some(new_hash.clone());
254        Ok(new_hash)
255    }
256
257    fn retract(&mut self, hash: &str, reason: &str) -> Result<()> {
258        let idx = *self
259            .by_hash
260            .get(hash)
261            .ok_or_else(|| Error::Substrate(format!("retract target {hash} not found")))?;
262        self.grains[idx].superseded_by = Some("retracted".to_string());
263        self.grains[idx]
264            .fields
265            .insert("verification_status".into(), json!("retracted"));
266        self.grains[idx]
267            .fields
268            .insert("retract_reason".into(), json!(reason));
269        Ok(())
270    }
271
272    fn execute_cal(&mut self, cal: &str) -> Result<Vec<Value>> {
273        let mut rows = Vec::new();
274        for line in cal.lines() {
275            let line = line.trim();
276            if line.is_empty() {
277                continue;
278            }
279            let (keyword, rest) = split_keyword(line);
280            match keyword.to_ascii_uppercase().as_str() {
281                "FORGET" => {
282                    let hash = rest.trim();
283                    if let Some(&idx) = self.by_hash.get(hash) {
284                        self.grains[idx].superseded_by = Some("forgotten".to_string());
285                    }
286                }
287                "RETRACT" => {
288                    let hash = rest.trim();
289                    self.retract(hash, "cal retract")?;
290                }
291                "ADD" => {
292                    let (grain_type, fields) = parse_type_and_json(rest)?;
293                    let spec = GrainSpec {
294                        grain_type,
295                        namespace: String::new(),
296                        fields,
297                    };
298                    let h = self.put_grain(&spec)?;
299                    rows.push(json!({ "hash": h }));
300                }
301                "SUPERSEDE" => {
302                    // SUPERSEDE <hash> WITH <type> {json}
303                    let (target, after_with) = rest.split_once(" WITH ").ok_or_else(|| {
304                        Error::CalUnsupported(format!("malformed SUPERSEDE: {line}"))
305                    })?;
306                    let (grain_type, fields) = parse_type_and_json(after_with)?;
307                    let spec = GrainSpec {
308                        grain_type,
309                        namespace: String::new(),
310                        fields,
311                    };
312                    let h = self.supersede(target.trim(), &spec, "cal supersede")?;
313                    rows.push(json!({ "hash": h }));
314                }
315                "DEFINE" => {
316                    let (key, _) = definition_key(line).ok_or_else(|| {
317                        Error::CalUnsupported(format!("malformed DEFINE: {line}"))
318                    })?;
319                    self.definitions.insert(key, line.to_string());
320                }
321                "DROP" => {
322                    if let Some((key, _)) = definition_key(line) {
323                        self.definitions.remove(&key);
324                    }
325                }
326                // Read verbs: no-ops in the reference substrate (no metric value).
327                "RECALL" | "ASSEMBLE" | "EXPLAIN" | "HISTORY" => {}
328                other => {
329                    return Err(Error::CalUnsupported(format!(
330                        "unknown statement {other:?}"
331                    )));
332                }
333            }
334        }
335        Ok(rows)
336    }
337
338    fn validate_cal(&self, cal: &str) -> Result<()> {
339        for line in cal.lines() {
340            let line = line.trim();
341            if line.is_empty() {
342                continue;
343            }
344            let (keyword, _) = split_keyword(line);
345            match keyword.to_ascii_uppercase().as_str() {
346                "ADD" | "SUPERSEDE" | "FORGET" | "RETRACT" | "RECALL" | "ASSEMBLE" | "EXPLAIN"
347                | "HISTORY" => {}
348                "DEFINE" | "DROP" => {
349                    if definition_key(line).is_none() {
350                        return Err(Error::CalUnsupported(format!(
351                            "malformed definition statement: {line}"
352                        )));
353                    }
354                }
355                other => {
356                    return Err(Error::CalUnsupported(format!(
357                        "unknown statement {other:?}"
358                    )))
359                }
360            }
361        }
362        Ok(())
363    }
364
365    /// The statement restoring the CURRENT definition, or a `DROP` when there
366    /// is none — so an apply can always record an inverse and a ROLLBACK
367    /// always undoes something.
368    fn definition_inverse(&self, statement: &str) -> Result<Option<String>> {
369        let Some((key, drop_stmt)) = definition_key(statement) else {
370            return Ok(None);
371        };
372        Ok(Some(
373            self.definitions.get(&key).cloned().unwrap_or(drop_stmt),
374        ))
375    }
376
377    fn load_state(&self) -> Result<Value> {
378        Ok(self.state.clone())
379    }
380
381    fn store_state(&mut self, state: &Value) -> Result<()> {
382        self.state = state.clone();
383        Ok(())
384    }
385}
386
387/// Parse a `DEFINE`/`DROP` statement into its registry key and the `DROP`
388/// that would remove it. `None` for anything else — this is a reader of two
389/// statement shapes, not a CAL parser.
390fn definition_key(line: &str) -> Option<(String, String)> {
391    let rest = {
392        let (kw, rest) = split_keyword(line.trim());
393        match kw.to_ascii_uppercase().as_str() {
394            "DEFINE" | "DROP" => rest,
395            _ => return None,
396        }
397    };
398    let (kind, rest) = split_keyword(rest);
399    let kind = kind.to_ascii_uppercase();
400    let (name, quoted) = match kind.as_str() {
401        "QUERY" => {
402            let rest = rest.trim().strip_prefix('"')?;
403            (rest.split('"').next()?, true)
404        }
405        "TEMPLATE" => (split_keyword(rest.trim()).0, false),
406        _ => return None,
407    };
408    if name.is_empty() {
409        return None;
410    }
411    let lower = kind.to_ascii_lowercase();
412    let drop_stmt = if quoted {
413        format!("DROP {kind} \"{name}\"")
414    } else {
415        format!("DROP {kind} {name}")
416    };
417    Some((format!("{lower}:{name}"), drop_stmt))
418}
419
420fn split_keyword(line: &str) -> (&str, &str) {
421    match line.split_once(char::is_whitespace) {
422        Some((k, rest)) => (k, rest.trim_start()),
423        None => (line, ""),
424    }
425}
426
427/// Parse `<type> {json}` → (type, fields).
428fn parse_type_and_json(s: &str) -> Result<(String, Map<String, Value>)> {
429    let brace = s
430        .find('{')
431        .ok_or_else(|| Error::CalUnsupported(format!("missing JSON object in {s:?}")))?;
432    let grain_type = s[..brace].trim().to_string();
433    if grain_type.is_empty() {
434        return Err(Error::CalUnsupported(format!(
435            "missing grain type in {s:?}"
436        )));
437    }
438    let value: Value = serde_json::from_str(s[brace..].trim())
439        .map_err(|e| Error::CalUnsupported(format!("bad JSON in {s:?}: {e}")))?;
440    let obj = value
441        .as_object()
442        .ok_or_else(|| Error::CalUnsupported(format!("JSON not an object in {s:?}")))?
443        .clone();
444    Ok((grain_type, obj))
445}
446
447#[cfg(test)]
448mod tests {
449    use super::*;
450
451    #[test]
452    fn forget_makes_grain_not_live() {
453        let mut sub = ReferenceSubstrate::new();
454        let h = sub
455            .put_grain(&GrainSpec::new("fact", "ns").with_field("subject", "x"))
456            .unwrap();
457        assert_eq!(
458            sub.grains_of_type("fact", None, ReadOpts::default())
459                .unwrap()
460                .len(),
461            1
462        );
463        sub.execute_cal(&format!("FORGET {h}")).unwrap();
464        assert!(sub
465            .grains_of_type("fact", None, ReadOpts::default())
466            .unwrap()
467            .is_empty());
468    }
469
470    #[test]
471    fn add_returns_hash_and_stores() {
472        let mut sub = ReferenceSubstrate::new();
473        let rows = sub
474            .execute_cal(r#"ADD fact {"subject":"acme","relation":"tier","object":"ent"}"#)
475            .unwrap();
476        assert_eq!(rows.len(), 1);
477        assert!(rows[0].get("hash").is_some());
478        assert_eq!(
479            sub.grains_of_type("fact", None, ReadOpts::default())
480                .unwrap()
481                .len(),
482            1
483        );
484    }
485
486    #[test]
487    fn validate_rejects_unknown_statement() {
488        let sub = ReferenceSubstrate::new();
489        assert!(sub.validate_cal("DROP TABLE").is_err());
490        assert!(sub.validate_cal("ADD fact {}").is_ok());
491    }
492
493    #[test]
494    fn state_round_trips() {
495        let mut sub = ReferenceSubstrate::new();
496        assert!(sub.load_state().unwrap().is_null());
497        sub.store_state(&json!({"k": 1})).unwrap();
498        assert_eq!(sub.load_state().unwrap(), json!({"k": 1}));
499    }
500}