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}
31
32impl ReferenceSubstrate {
33    pub fn new() -> Self {
34        ReferenceSubstrate {
35            state: Value::Null,
36            ..Default::default()
37        }
38    }
39
40    pub fn set_capabilities(&mut self, caps: Capabilities) {
41        self.caps = caps;
42    }
43
44    /// Register a fork (turns on the `forks` capability).
45    pub fn register_fork(&mut self, entity: &str, heads: &[&str]) {
46        self.caps.forks = true;
47        self.heads_index.insert(
48            entity.to_string(),
49            heads.iter().map(|s| s.to_string()).collect(),
50        );
51    }
52
53    /// Inject a telemetry snapshot (turns on the `telemetry` capability).
54    pub fn set_telemetry(&mut self, view: TelemetryView) {
55        self.caps.telemetry = true;
56        self.telemetry = Some(view);
57    }
58
59    /// Insert a fully-formed grain record; returns its assigned hash.
60    pub fn insert(&mut self, mut record: GrainRecord) -> String {
61        let hash = if record.hash.is_empty() {
62            self.mint_hash()
63        } else {
64            record.hash.clone()
65        };
66        record.hash = hash.clone();
67        let idx = self.grains.len();
68        self.by_hash.insert(hash.clone(), idx);
69        self.grains.push(record);
70        hash
71    }
72
73    fn mint_hash(&mut self) -> String {
74        let h = format!("ref-{:08}", self.next_id);
75        self.next_id += 1;
76        h
77    }
78
79    fn tick(&mut self) -> i64 {
80        self.clock += 1;
81        self.clock
82    }
83
84    fn record_from_spec(&mut self, spec: &GrainSpec) -> GrainRecord {
85        let created = self.tick();
86        let namespace = spec
87            .fields
88            .get("namespace")
89            .and_then(Value::as_str)
90            .unwrap_or(&spec.namespace)
91            .to_string();
92        let valid_to_ms = spec.fields.get("valid_to_ms").and_then(Value::as_i64);
93        GrainRecord {
94            hash: String::new(),
95            grain_type: spec.grain_type.clone(),
96            namespace,
97            created_at_ms: created,
98            valid_to_ms,
99            superseded_by: None,
100            fields: spec.fields.clone(),
101        }
102    }
103}
104
105impl SubstrateRead for ReferenceSubstrate {
106    fn capabilities(&self) -> Capabilities {
107        self.caps
108    }
109
110    fn grains_of_type(
111        &self,
112        grain_type: &str,
113        namespace: Option<&str>,
114        opts: ReadOpts,
115    ) -> Result<Vec<GrainRecord>> {
116        Ok(self
117            .grains
118            .iter()
119            .filter(|g| g.grain_type == grain_type)
120            .filter(|g| namespace.is_none_or(|ns| g.namespace == ns))
121            .filter(|g| !opts.live_only || g.is_live())
122            .filter(|g| opts.since_ms.is_none_or(|s| g.created_at_ms >= s))
123            .cloned()
124            .collect())
125    }
126
127    fn grain(&self, hash: &str) -> Result<Option<GrainRecord>> {
128        Ok(self.by_hash.get(hash).map(|&i| self.grains[i].clone()))
129    }
130
131    fn heads(&self, _namespace: Option<&str>) -> Result<Vec<HeadGroup>> {
132        if !self.caps.forks {
133            return Err(Error::CapabilityMissing("forks".into()));
134        }
135        let mut groups: Vec<HeadGroup> = self
136            .heads_index
137            .iter()
138            .map(|(entity, heads)| HeadGroup {
139                entity: entity.clone(),
140                heads: heads.clone(),
141            })
142            .collect();
143        groups.sort_by(|a, b| a.entity.cmp(&b.entity));
144        Ok(groups)
145    }
146
147    fn telemetry(&self, _namespace: Option<&str>) -> Result<Option<TelemetryView>> {
148        Ok(self.telemetry.clone())
149    }
150}
151
152impl OmsSubstrate for ReferenceSubstrate {
153    fn put_grain(&mut self, spec: &GrainSpec) -> Result<String> {
154        let record = self.record_from_spec(spec);
155        Ok(self.insert(record))
156    }
157
158    fn supersede(
159        &mut self,
160        target_hash: &str,
161        spec: &GrainSpec,
162        _justification: &str,
163    ) -> Result<String> {
164        let record = self.record_from_spec(spec);
165        let new_hash = self.insert(record);
166        let idx = *self
167            .by_hash
168            .get(target_hash)
169            .ok_or_else(|| Error::Substrate(format!("supersede target {target_hash} not found")))?;
170        self.grains[idx].superseded_by = Some(new_hash.clone());
171        Ok(new_hash)
172    }
173
174    fn retract(&mut self, hash: &str, reason: &str) -> Result<()> {
175        let idx = *self
176            .by_hash
177            .get(hash)
178            .ok_or_else(|| Error::Substrate(format!("retract target {hash} not found")))?;
179        self.grains[idx].superseded_by = Some("retracted".to_string());
180        self.grains[idx]
181            .fields
182            .insert("verification_status".into(), json!("retracted"));
183        self.grains[idx]
184            .fields
185            .insert("retract_reason".into(), json!(reason));
186        Ok(())
187    }
188
189    fn execute_cal(&mut self, cal: &str) -> Result<Vec<Value>> {
190        let mut rows = Vec::new();
191        for line in cal.lines() {
192            let line = line.trim();
193            if line.is_empty() {
194                continue;
195            }
196            let (keyword, rest) = split_keyword(line);
197            match keyword.to_ascii_uppercase().as_str() {
198                "FORGET" => {
199                    let hash = rest.trim();
200                    if let Some(&idx) = self.by_hash.get(hash) {
201                        self.grains[idx].superseded_by = Some("forgotten".to_string());
202                    }
203                }
204                "RETRACT" => {
205                    let hash = rest.trim();
206                    self.retract(hash, "cal retract")?;
207                }
208                "ADD" => {
209                    let (grain_type, fields) = parse_type_and_json(rest)?;
210                    let spec = GrainSpec {
211                        grain_type,
212                        namespace: String::new(),
213                        fields,
214                    };
215                    let h = self.put_grain(&spec)?;
216                    rows.push(json!({ "hash": h }));
217                }
218                "SUPERSEDE" => {
219                    // SUPERSEDE <hash> WITH <type> {json}
220                    let (target, after_with) = rest.split_once(" WITH ").ok_or_else(|| {
221                        Error::CalUnsupported(format!("malformed SUPERSEDE: {line}"))
222                    })?;
223                    let (grain_type, fields) = parse_type_and_json(after_with)?;
224                    let spec = GrainSpec {
225                        grain_type,
226                        namespace: String::new(),
227                        fields,
228                    };
229                    let h = self.supersede(target.trim(), &spec, "cal supersede")?;
230                    rows.push(json!({ "hash": h }));
231                }
232                // Read verbs: no-ops in the reference substrate (no metric value).
233                "RECALL" | "ASSEMBLE" | "EXPLAIN" | "HISTORY" => {}
234                other => {
235                    return Err(Error::CalUnsupported(format!(
236                        "unknown statement {other:?}"
237                    )));
238                }
239            }
240        }
241        Ok(rows)
242    }
243
244    fn validate_cal(&self, cal: &str) -> Result<()> {
245        for line in cal.lines() {
246            let line = line.trim();
247            if line.is_empty() {
248                continue;
249            }
250            let (keyword, _) = split_keyword(line);
251            match keyword.to_ascii_uppercase().as_str() {
252                "ADD" | "SUPERSEDE" | "FORGET" | "RETRACT" | "RECALL" | "ASSEMBLE" | "EXPLAIN"
253                | "HISTORY" => {}
254                other => {
255                    return Err(Error::CalUnsupported(format!(
256                        "unknown statement {other:?}"
257                    )))
258                }
259            }
260        }
261        Ok(())
262    }
263
264    fn load_state(&self) -> Result<Value> {
265        Ok(self.state.clone())
266    }
267
268    fn store_state(&mut self, state: &Value) -> Result<()> {
269        self.state = state.clone();
270        Ok(())
271    }
272}
273
274fn split_keyword(line: &str) -> (&str, &str) {
275    match line.split_once(char::is_whitespace) {
276        Some((k, rest)) => (k, rest.trim_start()),
277        None => (line, ""),
278    }
279}
280
281/// Parse `<type> {json}` → (type, fields).
282fn parse_type_and_json(s: &str) -> Result<(String, Map<String, Value>)> {
283    let brace = s
284        .find('{')
285        .ok_or_else(|| Error::CalUnsupported(format!("missing JSON object in {s:?}")))?;
286    let grain_type = s[..brace].trim().to_string();
287    if grain_type.is_empty() {
288        return Err(Error::CalUnsupported(format!(
289            "missing grain type in {s:?}"
290        )));
291    }
292    let value: Value = serde_json::from_str(s[brace..].trim())
293        .map_err(|e| Error::CalUnsupported(format!("bad JSON in {s:?}: {e}")))?;
294    let obj = value
295        .as_object()
296        .ok_or_else(|| Error::CalUnsupported(format!("JSON not an object in {s:?}")))?
297        .clone();
298    Ok((grain_type, obj))
299}
300
301#[cfg(test)]
302mod tests {
303    use super::*;
304
305    #[test]
306    fn forget_makes_grain_not_live() {
307        let mut sub = ReferenceSubstrate::new();
308        let h = sub
309            .put_grain(&GrainSpec::new("fact", "ns").with_field("subject", "x"))
310            .unwrap();
311        assert_eq!(
312            sub.grains_of_type("fact", None, ReadOpts::default())
313                .unwrap()
314                .len(),
315            1
316        );
317        sub.execute_cal(&format!("FORGET {h}")).unwrap();
318        assert!(sub
319            .grains_of_type("fact", None, ReadOpts::default())
320            .unwrap()
321            .is_empty());
322    }
323
324    #[test]
325    fn add_returns_hash_and_stores() {
326        let mut sub = ReferenceSubstrate::new();
327        let rows = sub
328            .execute_cal(r#"ADD fact {"subject":"acme","relation":"tier","object":"ent"}"#)
329            .unwrap();
330        assert_eq!(rows.len(), 1);
331        assert!(rows[0].get("hash").is_some());
332        assert_eq!(
333            sub.grains_of_type("fact", None, ReadOpts::default())
334                .unwrap()
335                .len(),
336            1
337        );
338    }
339
340    #[test]
341    fn validate_rejects_unknown_statement() {
342        let sub = ReferenceSubstrate::new();
343        assert!(sub.validate_cal("DROP TABLE").is_err());
344        assert!(sub.validate_cal("ADD fact {}").is_ok());
345    }
346
347    #[test]
348    fn state_round_trips() {
349        let mut sub = ReferenceSubstrate::new();
350        assert!(sub.load_state().unwrap().is_null());
351        sub.store_state(&json!({"k": 1})).unwrap();
352        assert_eq!(sub.load_state().unwrap(), json!({"k": 1}));
353    }
354}