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| 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 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 "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 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
387fn 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
427fn 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}