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 mock_embedder: bool,
24 plan_replay: Option<Value>,
26 state: Value,
27 next_id: u64,
28 clock: i64,
29 heads_index: HashMap<String, Vec<String>>,
31 telemetry: Option<TelemetryView>,
33 definitions: HashMap<String, String>,
37}
38
39impl ReferenceSubstrate {
40 pub fn new() -> Self {
41 ReferenceSubstrate {
42 state: Value::Null,
43 caps: Capabilities {
47 plans: true,
48 code: true,
49 ..Capabilities::default()
50 },
51 mock_embedder: false,
52 plan_replay: None,
53 ..Default::default()
54 }
55 }
56
57 pub fn set_mock_embedder(&mut self) {
63 self.mock_embedder = true;
64 self.caps.embeddings = true;
65 }
66
67 pub fn set_plan_replay(&mut self, report: Value) {
71 self.plan_replay = Some(report);
72 }
73
74 pub fn set_capabilities(&mut self, caps: Capabilities) {
75 self.caps = caps;
76 }
77
78 pub fn register_fork(&mut self, entity: &str, heads: &[&str]) {
80 self.caps.forks = true;
81 self.heads_index.insert(
82 entity.to_string(),
83 heads.iter().map(|s| s.to_string()).collect(),
84 );
85 }
86
87 pub fn set_telemetry(&mut self, view: TelemetryView) {
89 self.caps.telemetry = true;
90 self.telemetry = Some(view);
91 }
92
93 pub fn insert(&mut self, mut record: GrainRecord) -> String {
95 let hash = if record.hash.is_empty() {
96 self.mint_hash()
97 } else {
98 record.hash.clone()
99 };
100 record.hash = hash.clone();
101 let idx = self.grains.len();
102 self.by_hash.insert(hash.clone(), idx);
103 self.grains.push(record);
104 hash
105 }
106
107 fn mint_hash(&mut self) -> String {
108 let h = format!("ref-{:08}", self.next_id);
109 self.next_id += 1;
110 h
111 }
112
113 fn tick(&mut self) -> i64 {
114 self.clock += 1;
115 self.clock
116 }
117
118 fn record_from_spec(&mut self, spec: &GrainSpec) -> GrainRecord {
119 let created = self.tick();
120 let namespace = spec
121 .fields
122 .get("namespace")
123 .and_then(Value::as_str)
124 .unwrap_or(&spec.namespace)
125 .to_string();
126 let valid_to_ms = spec.fields.get("valid_to_ms").and_then(Value::as_i64);
127 GrainRecord {
128 hash: String::new(),
129 grain_type: spec.grain_type.clone(),
130 namespace,
131 created_at_ms: created,
132 valid_to_ms,
133 superseded_by: None,
134 fields: spec.fields.clone(),
135 }
136 }
137}
138
139impl SubstrateRead for ReferenceSubstrate {
140 fn capabilities(&self) -> Capabilities {
141 self.caps
142 }
143
144 fn embed(&self, text: &str) -> Result<Option<Vec<f32>>> {
145 if !self.mock_embedder {
146 return Ok(None);
147 }
148 Ok(Some(mock_embed(text)))
149 }
150
151 fn plan_replay(&self, _incumbent: &str, _candidate: &Value) -> Result<Option<Value>> {
152 Ok(self.plan_replay.clone())
153 }
154
155 fn validate_plan(&self, workflow: &Value) -> Result<()> {
161 let bad = |m: &str| Err(Error::Substrate(m.to_string()));
162 let Some(fields) = workflow.as_object() else {
163 return bad("a workflow body must be a JSON object");
164 };
165 let Some(Value::Array(nodes)) = fields.get("nodes") else {
166 return bad("a workflow needs a 'nodes' array");
167 };
168 let mut seen = std::collections::BTreeSet::new();
169 for n in nodes {
170 match n.as_str() {
171 Some(id) if seen.insert(id) => {}
172 Some(id) => return bad(&format!("duplicate node {id:?}")),
173 None => return bad("every workflow node must be a string"),
174 }
175 }
176 if seen.is_empty() {
177 return bad("a workflow needs at least one node");
178 }
179 if let Some(v) = fields.get("edges") {
180 let Some(edges) = v.as_array() else {
181 return bad("workflow 'edges' must be an array");
182 };
183 for e in edges {
184 for end in ["src", "dst"] {
185 match e.get(end).and_then(Value::as_str) {
186 Some(id) if seen.contains(id) => {}
187 Some(id) => return bad(&format!("edge {end} {id:?} is not a node")),
188 None => return bad(&format!("every edge needs a '{end}'")),
189 }
190 }
191 }
192 }
193 Ok(())
194 }
195
196 fn tool_evalset(&self, tool: &str) -> Result<Option<String>> {
198 Ok(self
199 .grains
200 .iter()
201 .filter(|g| {
202 g.grain_type == crate::model::grain_type::TOOL
203 && g.is_live()
204 && g.str_field("kind") == Some("definition")
205 && g.tool_name() == Some(tool)
206 })
207 .find_map(|g| {
208 g.str_field("evalset_hash")
209 .filter(|h| !h.trim().is_empty())
210 .map(str::to_string)
211 }))
212 }
213
214 fn grains_of_type(
215 &self,
216 grain_type: &str,
217 namespace: Option<&str>,
218 opts: ReadOpts,
219 ) -> Result<Vec<GrainRecord>> {
220 Ok(self
221 .grains
222 .iter()
223 .filter(|g| g.grain_type == grain_type)
224 .filter(|g| namespace.is_none_or(|ns| g.namespace == ns))
225 .filter(|g| namespace.is_some() || !g.namespace.starts_with("agent:"))
237 .filter(|g| !opts.live_only || g.is_live())
238 .filter(|g| opts.since_ms.is_none_or(|s| g.created_at_ms >= s))
239 .cloned()
240 .collect())
241 }
242
243 fn grain(&self, hash: &str) -> Result<Option<GrainRecord>> {
244 Ok(self.by_hash.get(hash).map(|&i| self.grains[i].clone()))
245 }
246
247 fn heads(&self, _namespace: Option<&str>) -> Result<Vec<HeadGroup>> {
248 if !self.caps.forks {
249 return Err(Error::CapabilityMissing("forks".into()));
250 }
251 let mut groups: Vec<HeadGroup> = self
252 .heads_index
253 .iter()
254 .map(|(entity, heads)| HeadGroup {
255 entity: entity.clone(),
256 heads: heads.clone(),
257 })
258 .collect();
259 groups.sort_by(|a, b| a.entity.cmp(&b.entity));
260 Ok(groups)
261 }
262
263 fn telemetry(&self, _namespace: Option<&str>) -> Result<Option<TelemetryView>> {
264 Ok(self.telemetry.clone())
265 }
266}
267
268impl OmsSubstrate for ReferenceSubstrate {
269 fn put_grain(&mut self, spec: &GrainSpec) -> Result<String> {
270 let record = self.record_from_spec(spec);
271 Ok(self.insert(record))
272 }
273
274 fn supersede(
275 &mut self,
276 target_hash: &str,
277 spec: &GrainSpec,
278 _justification: &str,
279 ) -> Result<String> {
280 let record = self.record_from_spec(spec);
281 let new_hash = self.insert(record);
282 let idx = *self
283 .by_hash
284 .get(target_hash)
285 .ok_or_else(|| Error::Substrate(format!("supersede target {target_hash} not found")))?;
286 self.grains[idx].superseded_by = Some(new_hash.clone());
287 Ok(new_hash)
288 }
289
290 fn retract(&mut self, hash: &str, reason: &str) -> Result<()> {
291 let idx = *self
292 .by_hash
293 .get(hash)
294 .ok_or_else(|| Error::Substrate(format!("retract target {hash} not found")))?;
295 self.grains[idx].superseded_by = Some("retracted".to_string());
296 self.grains[idx]
297 .fields
298 .insert("verification_status".into(), json!("retracted"));
299 self.grains[idx]
300 .fields
301 .insert("retract_reason".into(), json!(reason));
302 for g in self.grains.iter_mut() {
306 if g.superseded_by.as_deref() == Some(hash) {
307 g.superseded_by = None;
308 }
309 }
310 Ok(())
311 }
312
313 fn execute_cal(&mut self, cal: &str) -> Result<Vec<Value>> {
314 let mut rows = Vec::new();
315 for line in cal.lines() {
316 let line = line.trim();
317 if line.is_empty() {
318 continue;
319 }
320 let (keyword, rest) = split_keyword(line);
321 match keyword.to_ascii_uppercase().as_str() {
322 "FORGET" => {
323 let hash = rest.trim();
324 if let Some(&idx) = self.by_hash.get(hash) {
325 self.grains[idx].superseded_by = Some("forgotten".to_string());
326 }
327 }
328 "RETRACT" => {
329 let hash = rest.trim();
330 self.retract(hash, "cal retract")?;
331 }
332 "ADD" => {
333 let (grain_type, fields) = parse_type_and_json(rest)?;
334 let spec = GrainSpec {
335 grain_type,
336 namespace: String::new(),
337 fields,
338 };
339 let h = self.put_grain(&spec)?;
340 rows.push(json!({ "hash": h }));
341 }
342 "SUPERSEDE" => {
343 let (target, after_with) = rest.split_once(" WITH ").ok_or_else(|| {
345 Error::CalUnsupported(format!("malformed SUPERSEDE: {line}"))
346 })?;
347 let (grain_type, fields) = parse_type_and_json(after_with)?;
348 let spec = GrainSpec {
349 grain_type,
350 namespace: String::new(),
351 fields,
352 };
353 let h = self.supersede(target.trim(), &spec, "cal supersede")?;
354 rows.push(json!({ "hash": h }));
355 }
356 "DEFINE" => {
357 let (key, _) = definition_key(line).ok_or_else(|| {
358 Error::CalUnsupported(format!("malformed DEFINE: {line}"))
359 })?;
360 self.definitions.insert(key, line.to_string());
361 }
362 "DROP" => {
363 if let Some((key, _)) = definition_key(line) {
364 self.definitions.remove(&key);
365 }
366 }
367 "RECALL" | "ASSEMBLE" | "EXPLAIN" | "HISTORY" => {}
369 other => {
370 return Err(Error::CalUnsupported(format!(
371 "unknown statement {other:?}"
372 )));
373 }
374 }
375 }
376 Ok(rows)
377 }
378
379 fn validate_cal(&self, cal: &str) -> Result<()> {
380 for line in cal.lines() {
381 let line = line.trim();
382 if line.is_empty() {
383 continue;
384 }
385 let (keyword, _) = split_keyword(line);
386 match keyword.to_ascii_uppercase().as_str() {
387 "ADD" | "SUPERSEDE" | "FORGET" | "RETRACT" | "RECALL" | "ASSEMBLE" | "EXPLAIN"
388 | "HISTORY" => {}
389 "DEFINE" | "DROP" => {
390 if definition_key(line).is_none() {
391 return Err(Error::CalUnsupported(format!(
392 "malformed definition statement: {line}"
393 )));
394 }
395 }
396 other => {
397 return Err(Error::CalUnsupported(format!(
398 "unknown statement {other:?}"
399 )))
400 }
401 }
402 }
403 Ok(())
404 }
405
406 fn definition_inverse(&self, statement: &str) -> Result<Option<String>> {
410 let Some((key, drop_stmt)) = definition_key(statement) else {
411 return Ok(None);
412 };
413 Ok(Some(
414 self.definitions.get(&key).cloned().unwrap_or(drop_stmt),
415 ))
416 }
417
418 fn load_state(&self) -> Result<Value> {
419 Ok(self.state.clone())
420 }
421
422 fn store_state(&mut self, state: &Value) -> Result<()> {
423 self.state = state.clone();
424 Ok(())
425 }
426}
427
428fn definition_key(line: &str) -> Option<(String, String)> {
432 let rest = {
433 let (kw, rest) = split_keyword(line.trim());
434 match kw.to_ascii_uppercase().as_str() {
435 "DEFINE" | "DROP" => rest,
436 _ => return None,
437 }
438 };
439 let (kind, rest) = split_keyword(rest);
440 let kind = kind.to_ascii_uppercase();
441 let (name, quoted) = match kind.as_str() {
442 "QUERY" => {
443 let rest = rest.trim().strip_prefix('"')?;
444 (rest.split('"').next()?, true)
445 }
446 "TEMPLATE" => (split_keyword(rest.trim()).0, false),
447 _ => return None,
448 };
449 if name.is_empty() {
450 return None;
451 }
452 let lower = kind.to_ascii_lowercase();
453 let drop_stmt = if quoted {
454 format!("DROP {kind} \"{name}\"")
455 } else {
456 format!("DROP {kind} {name}")
457 };
458 Some((format!("{lower}:{name}"), drop_stmt))
459}
460
461fn split_keyword(line: &str) -> (&str, &str) {
462 match line.split_once(char::is_whitespace) {
463 Some((k, rest)) => (k, rest.trim_start()),
464 None => (line, ""),
465 }
466}
467
468fn parse_type_and_json(s: &str) -> Result<(String, Map<String, Value>)> {
470 let brace = s
471 .find('{')
472 .ok_or_else(|| Error::CalUnsupported(format!("missing JSON object in {s:?}")))?;
473 let grain_type = s[..brace].trim().to_string();
474 if grain_type.is_empty() {
475 return Err(Error::CalUnsupported(format!(
476 "missing grain type in {s:?}"
477 )));
478 }
479 let value: Value = serde_json::from_str(s[brace..].trim())
480 .map_err(|e| Error::CalUnsupported(format!("bad JSON in {s:?}: {e}")))?;
481 let obj = value
482 .as_object()
483 .ok_or_else(|| Error::CalUnsupported(format!("JSON not an object in {s:?}")))?
484 .clone();
485 Ok((grain_type, obj))
486}
487
488fn mock_embed(text: &str) -> Vec<f32> {
493 const SYNONYMS: &[(&str, &str)] = &[
494 ("write", "record"), ("note", "record"), ("log", "record"), ("capture", "record"),
495 ("down", ""), ("supplier", "vendor"), ("seller", "vendor"), ("merchant", "vendor"),
496 ("always", ""), ("every", "each"), ("all", "each"), ("a", ""), ("an", ""), ("the", ""),
497 ("on", ""), ("of", ""), ("s", ""), ("for", ""), ("to", ""), ("and", ""), ("with", ""),
498 ("before", "prior"), ("ahead", "prior"), ("answering", "answer"), ("answers", "answer"),
499 ("confirm", "check"), ("verify", "check"), ("current", "present"), ("latest", "present"),
500 ("city", "location"), ("town", "location"), ("place", "location"),
501 ("invoice", "bill"), ("receipt", "bill"), ("number", "id"), ("identifier", "id"),
502 ("exactly", "verbatim"), ("printed", "shown"),
503 ];
504 let mut v = vec![0f32; 64];
505 for raw in text.to_lowercase().split(|c: char| !c.is_alphanumeric()) {
506 if raw.is_empty() {
507 continue;
508 }
509 let tok = SYNONYMS.iter().find(|(from, _)| *from == raw).map(|(_, to)| *to).unwrap_or(raw);
510 if tok.is_empty() {
511 continue;
512 }
513 let mut h: u64 = 0xcbf29ce484222325;
515 for b in tok.bytes() {
516 h ^= b as u64;
517 h = h.wrapping_mul(0x100000001b3);
518 }
519 v[(h % 64) as usize] += 1.0;
520 }
521 v
522}
523
524#[cfg(test)]
525mod tests {
526 use super::*;
527
528 #[test]
529 fn forget_makes_grain_not_live() {
530 let mut sub = ReferenceSubstrate::new();
531 let h = sub
532 .put_grain(&GrainSpec::new("fact", "ns").with_field("subject", "x"))
533 .unwrap();
534 assert_eq!(
535 sub.grains_of_type("fact", None, ReadOpts::default())
536 .unwrap()
537 .len(),
538 1
539 );
540 sub.execute_cal(&format!("FORGET {h}")).unwrap();
541 assert!(sub
542 .grains_of_type("fact", None, ReadOpts::default())
543 .unwrap()
544 .is_empty());
545 }
546
547 #[test]
548 fn add_returns_hash_and_stores() {
549 let mut sub = ReferenceSubstrate::new();
550 let rows = sub
551 .execute_cal(r#"ADD fact {"subject":"acme","relation":"tier","object":"ent"}"#)
552 .unwrap();
553 assert_eq!(rows.len(), 1);
554 assert!(rows[0].get("hash").is_some());
555 assert_eq!(
556 sub.grains_of_type("fact", None, ReadOpts::default())
557 .unwrap()
558 .len(),
559 1
560 );
561 }
562
563 #[test]
564 fn validate_rejects_unknown_statement() {
565 let sub = ReferenceSubstrate::new();
566 assert!(sub.validate_cal("DROP TABLE").is_err());
567 assert!(sub.validate_cal("ADD fact {}").is_ok());
568 }
569
570 #[test]
571 fn state_round_trips() {
572 let mut sub = ReferenceSubstrate::new();
573 assert!(sub.load_state().unwrap().is_null());
574 sub.store_state(&json!({"k": 1})).unwrap();
575 assert_eq!(sub.load_state().unwrap(), json!({"k": 1}));
576 }
577}