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            .filter(|g| !opts.live_only || g.is_live())
193            .filter(|g| opts.since_ms.is_none_or(|s| g.created_at_ms >= s))
194            .cloned()
195            .collect())
196    }
197
198    fn grain(&self, hash: &str) -> Result<Option<GrainRecord>> {
199        Ok(self.by_hash.get(hash).map(|&i| self.grains[i].clone()))
200    }
201
202    fn heads(&self, _namespace: Option<&str>) -> Result<Vec<HeadGroup>> {
203        if !self.caps.forks {
204            return Err(Error::CapabilityMissing("forks".into()));
205        }
206        let mut groups: Vec<HeadGroup> = self
207            .heads_index
208            .iter()
209            .map(|(entity, heads)| HeadGroup {
210                entity: entity.clone(),
211                heads: heads.clone(),
212            })
213            .collect();
214        groups.sort_by(|a, b| a.entity.cmp(&b.entity));
215        Ok(groups)
216    }
217
218    fn telemetry(&self, _namespace: Option<&str>) -> Result<Option<TelemetryView>> {
219        Ok(self.telemetry.clone())
220    }
221}
222
223impl OmsSubstrate for ReferenceSubstrate {
224    fn put_grain(&mut self, spec: &GrainSpec) -> Result<String> {
225        let record = self.record_from_spec(spec);
226        Ok(self.insert(record))
227    }
228
229    fn supersede(
230        &mut self,
231        target_hash: &str,
232        spec: &GrainSpec,
233        _justification: &str,
234    ) -> Result<String> {
235        let record = self.record_from_spec(spec);
236        let new_hash = self.insert(record);
237        let idx = *self
238            .by_hash
239            .get(target_hash)
240            .ok_or_else(|| Error::Substrate(format!("supersede target {target_hash} not found")))?;
241        self.grains[idx].superseded_by = Some(new_hash.clone());
242        Ok(new_hash)
243    }
244
245    fn retract(&mut self, hash: &str, reason: &str) -> Result<()> {
246        let idx = *self
247            .by_hash
248            .get(hash)
249            .ok_or_else(|| Error::Substrate(format!("retract target {hash} not found")))?;
250        self.grains[idx].superseded_by = Some("retracted".to_string());
251        self.grains[idx]
252            .fields
253            .insert("verification_status".into(), json!("retracted"));
254        self.grains[idx]
255            .fields
256            .insert("retract_reason".into(), json!(reason));
257        Ok(())
258    }
259
260    fn execute_cal(&mut self, cal: &str) -> Result<Vec<Value>> {
261        let mut rows = Vec::new();
262        for line in cal.lines() {
263            let line = line.trim();
264            if line.is_empty() {
265                continue;
266            }
267            let (keyword, rest) = split_keyword(line);
268            match keyword.to_ascii_uppercase().as_str() {
269                "FORGET" => {
270                    let hash = rest.trim();
271                    if let Some(&idx) = self.by_hash.get(hash) {
272                        self.grains[idx].superseded_by = Some("forgotten".to_string());
273                    }
274                }
275                "RETRACT" => {
276                    let hash = rest.trim();
277                    self.retract(hash, "cal retract")?;
278                }
279                "ADD" => {
280                    let (grain_type, fields) = parse_type_and_json(rest)?;
281                    let spec = GrainSpec {
282                        grain_type,
283                        namespace: String::new(),
284                        fields,
285                    };
286                    let h = self.put_grain(&spec)?;
287                    rows.push(json!({ "hash": h }));
288                }
289                "SUPERSEDE" => {
290                    // SUPERSEDE <hash> WITH <type> {json}
291                    let (target, after_with) = rest.split_once(" WITH ").ok_or_else(|| {
292                        Error::CalUnsupported(format!("malformed SUPERSEDE: {line}"))
293                    })?;
294                    let (grain_type, fields) = parse_type_and_json(after_with)?;
295                    let spec = GrainSpec {
296                        grain_type,
297                        namespace: String::new(),
298                        fields,
299                    };
300                    let h = self.supersede(target.trim(), &spec, "cal supersede")?;
301                    rows.push(json!({ "hash": h }));
302                }
303                "DEFINE" => {
304                    let (key, _) = definition_key(line).ok_or_else(|| {
305                        Error::CalUnsupported(format!("malformed DEFINE: {line}"))
306                    })?;
307                    self.definitions.insert(key, line.to_string());
308                }
309                "DROP" => {
310                    if let Some((key, _)) = definition_key(line) {
311                        self.definitions.remove(&key);
312                    }
313                }
314                // Read verbs: no-ops in the reference substrate (no metric value).
315                "RECALL" | "ASSEMBLE" | "EXPLAIN" | "HISTORY" => {}
316                other => {
317                    return Err(Error::CalUnsupported(format!(
318                        "unknown statement {other:?}"
319                    )));
320                }
321            }
322        }
323        Ok(rows)
324    }
325
326    fn validate_cal(&self, cal: &str) -> Result<()> {
327        for line in cal.lines() {
328            let line = line.trim();
329            if line.is_empty() {
330                continue;
331            }
332            let (keyword, _) = split_keyword(line);
333            match keyword.to_ascii_uppercase().as_str() {
334                "ADD" | "SUPERSEDE" | "FORGET" | "RETRACT" | "RECALL" | "ASSEMBLE" | "EXPLAIN"
335                | "HISTORY" => {}
336                "DEFINE" | "DROP" => {
337                    if definition_key(line).is_none() {
338                        return Err(Error::CalUnsupported(format!(
339                            "malformed definition statement: {line}"
340                        )));
341                    }
342                }
343                other => {
344                    return Err(Error::CalUnsupported(format!(
345                        "unknown statement {other:?}"
346                    )))
347                }
348            }
349        }
350        Ok(())
351    }
352
353    /// The statement restoring the CURRENT definition, or a `DROP` when there
354    /// is none — so an apply can always record an inverse and a ROLLBACK
355    /// always undoes something.
356    fn definition_inverse(&self, statement: &str) -> Result<Option<String>> {
357        let Some((key, drop_stmt)) = definition_key(statement) else {
358            return Ok(None);
359        };
360        Ok(Some(
361            self.definitions.get(&key).cloned().unwrap_or(drop_stmt),
362        ))
363    }
364
365    fn load_state(&self) -> Result<Value> {
366        Ok(self.state.clone())
367    }
368
369    fn store_state(&mut self, state: &Value) -> Result<()> {
370        self.state = state.clone();
371        Ok(())
372    }
373}
374
375/// Parse a `DEFINE`/`DROP` statement into its registry key and the `DROP`
376/// that would remove it. `None` for anything else — this is a reader of two
377/// statement shapes, not a CAL parser.
378fn definition_key(line: &str) -> Option<(String, String)> {
379    let rest = {
380        let (kw, rest) = split_keyword(line.trim());
381        match kw.to_ascii_uppercase().as_str() {
382            "DEFINE" | "DROP" => rest,
383            _ => return None,
384        }
385    };
386    let (kind, rest) = split_keyword(rest);
387    let kind = kind.to_ascii_uppercase();
388    let (name, quoted) = match kind.as_str() {
389        "QUERY" => {
390            let rest = rest.trim().strip_prefix('"')?;
391            (rest.split('"').next()?, true)
392        }
393        "TEMPLATE" => (split_keyword(rest.trim()).0, false),
394        _ => return None,
395    };
396    if name.is_empty() {
397        return None;
398    }
399    let lower = kind.to_ascii_lowercase();
400    let drop_stmt = if quoted {
401        format!("DROP {kind} \"{name}\"")
402    } else {
403        format!("DROP {kind} {name}")
404    };
405    Some((format!("{lower}:{name}"), drop_stmt))
406}
407
408fn split_keyword(line: &str) -> (&str, &str) {
409    match line.split_once(char::is_whitespace) {
410        Some((k, rest)) => (k, rest.trim_start()),
411        None => (line, ""),
412    }
413}
414
415/// Parse `<type> {json}` → (type, fields).
416fn parse_type_and_json(s: &str) -> Result<(String, Map<String, Value>)> {
417    let brace = s
418        .find('{')
419        .ok_or_else(|| Error::CalUnsupported(format!("missing JSON object in {s:?}")))?;
420    let grain_type = s[..brace].trim().to_string();
421    if grain_type.is_empty() {
422        return Err(Error::CalUnsupported(format!(
423            "missing grain type in {s:?}"
424        )));
425    }
426    let value: Value = serde_json::from_str(s[brace..].trim())
427        .map_err(|e| Error::CalUnsupported(format!("bad JSON in {s:?}: {e}")))?;
428    let obj = value
429        .as_object()
430        .ok_or_else(|| Error::CalUnsupported(format!("JSON not an object in {s:?}")))?
431        .clone();
432    Ok((grain_type, obj))
433}
434
435#[cfg(test)]
436mod tests {
437    use super::*;
438
439    #[test]
440    fn forget_makes_grain_not_live() {
441        let mut sub = ReferenceSubstrate::new();
442        let h = sub
443            .put_grain(&GrainSpec::new("fact", "ns").with_field("subject", "x"))
444            .unwrap();
445        assert_eq!(
446            sub.grains_of_type("fact", None, ReadOpts::default())
447                .unwrap()
448                .len(),
449            1
450        );
451        sub.execute_cal(&format!("FORGET {h}")).unwrap();
452        assert!(sub
453            .grains_of_type("fact", None, ReadOpts::default())
454            .unwrap()
455            .is_empty());
456    }
457
458    #[test]
459    fn add_returns_hash_and_stores() {
460        let mut sub = ReferenceSubstrate::new();
461        let rows = sub
462            .execute_cal(r#"ADD fact {"subject":"acme","relation":"tier","object":"ent"}"#)
463            .unwrap();
464        assert_eq!(rows.len(), 1);
465        assert!(rows[0].get("hash").is_some());
466        assert_eq!(
467            sub.grains_of_type("fact", None, ReadOpts::default())
468                .unwrap()
469                .len(),
470            1
471        );
472    }
473
474    #[test]
475    fn validate_rejects_unknown_statement() {
476        let sub = ReferenceSubstrate::new();
477        assert!(sub.validate_cal("DROP TABLE").is_err());
478        assert!(sub.validate_cal("ADD fact {}").is_ok());
479    }
480
481    #[test]
482    fn state_round_trips() {
483        let mut sub = ReferenceSubstrate::new();
484        assert!(sub.load_state().unwrap().is_null());
485        sub.store_state(&json!({"k": 1})).unwrap();
486        assert_eq!(sub.load_state().unwrap(), json!({"k": 1}));
487    }
488}