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 definitions: HashMap<String, String>,
34}
35
36impl ReferenceSubstrate {
37 pub fn new() -> Self {
38 ReferenceSubstrate {
39 state: Value::Null,
40 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 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 pub fn set_telemetry(&mut self, view: TelemetryView) {
67 self.caps.telemetry = true;
68 self.telemetry = Some(view);
69 }
70
71 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 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 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 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 "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 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
375fn 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
415fn 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}