1use 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 heads_index: HashMap<String, Vec<String>>,
28 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 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 pub fn set_telemetry(&mut self, view: TelemetryView) {
55 self.caps.telemetry = true;
56 self.telemetry = Some(view);
57 }
58
59 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 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 "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
281fn 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}