1use serde::{Deserialize, Serialize};
2use std::collections::BTreeMap;
3
4use crate::clock::LamportClock;
5use crate::ontology::{Ontology, OntologyExtension};
6
7#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
9#[serde(untagged)]
10pub enum Value {
11 Null,
12 Bool(bool),
13 Int(i64),
14 Float(f64),
15 String(String),
16 List(Vec<Value>),
17 Map(BTreeMap<String, Value>),
18}
19
20#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
25#[serde(tag = "op")]
26pub enum GraphOp {
27 #[serde(rename = "define_ontology")]
30 DefineOntology { ontology: Ontology },
31 #[serde(rename = "add_node")]
32 AddNode {
33 node_id: String,
34 node_type: String,
35 #[serde(default)]
36 subtype: Option<String>,
37 label: String,
38 #[serde(default)]
39 properties: BTreeMap<String, Value>,
40 },
41 #[serde(rename = "add_edge")]
42 AddEdge {
43 edge_id: String,
44 edge_type: String,
45 source_id: String,
46 target_id: String,
47 #[serde(default)]
48 properties: BTreeMap<String, Value>,
49 },
50 #[serde(rename = "update_property")]
51 UpdateProperty {
52 entity_id: String,
53 key: String,
54 value: Value,
55 },
56 #[serde(rename = "remove_node")]
57 RemoveNode { node_id: String },
58 #[serde(rename = "remove_edge")]
59 RemoveEdge { edge_id: String },
60 #[serde(rename = "extend_ontology")]
62 ExtendOntology { extension: OntologyExtension },
63 #[serde(rename = "define_lens")]
67 DefineLens { transforms: Vec<u8> },
68 #[serde(rename = "checkpoint")]
72 Checkpoint {
73 ops: Vec<GraphOp>,
75 #[serde(default)]
78 op_clocks: Vec<(u64, u32)>,
79 compacted_at_physical_ms: u64,
81 compacted_at_logical: u32,
83 },
84}
85
86pub type Hash = [u8; 32];
88
89#[derive(Debug, Clone, PartialEq, Serialize)]
94pub struct Entry {
95 pub hash: Hash,
97 pub payload: GraphOp,
99 pub next: Vec<Hash>,
101 #[serde(default)]
104 pub refs: Vec<Hash>,
105 pub clock: LamportClock,
107 pub author: String,
109 #[serde(default)]
111 pub signature: Option<Vec<u8>>,
112}
113
114impl<'de> Deserialize<'de> for Entry {
128 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
129 where
130 D: serde::Deserializer<'de>,
131 {
132 struct EntryVisitor;
133
134 impl<'de> serde::de::Visitor<'de> for EntryVisitor {
135 type Value = Entry;
136
137 fn expecting(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
138 f.write_str("an Entry as a 7-element sequence (or legacy 8-element)")
139 }
140
141 fn visit_seq<A>(self, mut seq: A) -> Result<Entry, A::Error>
142 where
143 A: serde::de::SeqAccess<'de>,
144 {
145 use serde::de::Error as _;
146 let hash = seq
147 .next_element()?
148 .ok_or_else(|| A::Error::invalid_length(0, &"8 fields"))?;
149 let payload = seq
150 .next_element()?
151 .ok_or_else(|| A::Error::invalid_length(1, &"8 fields"))?;
152 let next = seq
153 .next_element()?
154 .ok_or_else(|| A::Error::invalid_length(2, &"8 fields"))?;
155 let refs = seq.next_element()?.unwrap_or_default();
156 let clock = seq
157 .next_element()?
158 .ok_or_else(|| A::Error::invalid_length(4, &"8 fields"))?;
159 let author = seq
160 .next_element()?
161 .ok_or_else(|| A::Error::invalid_length(5, &"8 fields"))?;
162 let signature = seq.next_element()?.unwrap_or(None);
163 let _legacy: Option<Option<Hash>> = seq.next_element()?;
166
167 Ok(Entry {
168 hash,
169 payload,
170 next,
171 refs,
172 clock,
173 author,
174 signature,
175 })
176 }
177 }
178
179 deserializer.deserialize_seq(EntryVisitor)
180 }
181}
182
183#[derive(Serialize)]
186struct SignableContent<'a> {
187 payload: &'a GraphOp,
188 next: &'a Vec<Hash>,
189 refs: &'a Vec<Hash>,
190 clock: &'a LamportClock,
191 author: &'a str,
192}
193
194impl Entry {
195 pub fn new(
197 payload: GraphOp,
198 next: Vec<Hash>,
199 refs: Vec<Hash>,
200 clock: LamportClock,
201 author: impl Into<String>,
202 ) -> Self {
203 let author = author.into();
204 let hash = Self::compute_hash(&payload, &next, &refs, &clock, &author);
205 Self {
206 hash,
207 payload,
208 next,
209 refs,
210 clock,
211 author,
212 signature: None,
213 }
214 }
215
216 #[cfg(feature = "signing")]
218 pub fn new_signed(
219 payload: GraphOp,
220 next: Vec<Hash>,
221 refs: Vec<Hash>,
222 clock: LamportClock,
223 author: impl Into<String>,
224 signing_key: &ed25519_dalek::SigningKey,
225 ) -> Self {
226 use ed25519_dalek::Signer;
227 let author = author.into();
228 let hash = Self::compute_hash(&payload, &next, &refs, &clock, &author);
229 let sig = signing_key.sign(&hash);
230 Self {
231 hash,
232 payload,
233 next,
234 refs,
235 clock,
236 author,
237 signature: Some(sig.to_bytes().to_vec()),
238 }
239 }
240
241 #[cfg(feature = "signing")]
245 pub fn verify_signature(&self, public_key: &ed25519_dalek::VerifyingKey) -> bool {
246 use ed25519_dalek::Verifier;
247 match &self.signature {
248 Some(sig_bytes) => {
249 if sig_bytes.len() != 64 {
250 return false;
251 }
252 let mut sig_array = [0u8; 64];
253 sig_array.copy_from_slice(sig_bytes);
254 let sig = ed25519_dalek::Signature::from_bytes(&sig_array);
255 public_key.verify(&self.hash, &sig).is_ok()
256 }
257 None => true, }
259 }
260
261 pub fn is_signed(&self) -> bool {
263 self.signature.is_some()
264 }
265
266 fn compute_hash(
268 payload: &GraphOp,
269 next: &Vec<Hash>,
270 refs: &Vec<Hash>,
271 clock: &LamportClock,
272 author: &str,
273 ) -> Hash {
274 let signable = SignableContent {
275 payload,
276 next,
277 refs,
278 clock,
279 author,
280 };
281 let bytes = rmp_serde::to_vec(&signable).expect("serialization should not fail");
284 *blake3::hash(&bytes).as_bytes()
285 }
286
287 pub fn verify_hash(&self) -> bool {
289 let computed = Self::compute_hash(
290 &self.payload,
291 &self.next,
292 &self.refs,
293 &self.clock,
294 &self.author,
295 );
296 self.hash == computed
297 }
298
299 pub fn to_bytes(&self) -> Vec<u8> {
305 rmp_serde::to_vec(self).expect("entry serialization should not fail")
306 }
307
308 pub fn from_bytes(bytes: &[u8]) -> Result<Self, rmp_serde::decode::Error> {
310 rmp_serde::from_slice(bytes)
311 }
312
313 pub fn hash_hex(&self) -> String {
315 hex::encode(self.hash)
316 }
317}
318
319pub fn hash_hex(hash: &Hash) -> String {
321 hex::encode(hash)
322}
323
324#[cfg(test)]
325mod tests {
326 use super::*;
327 use crate::ontology::{EdgeTypeDef, NodeTypeDef, PropertyDef, ValueType};
328
329 fn sample_ontology() -> Ontology {
330 Ontology {
331 node_types: BTreeMap::from([
332 (
333 "entity".into(),
334 NodeTypeDef {
335 description: None,
336 properties: BTreeMap::from([
337 (
338 "ip".into(),
339 PropertyDef {
340 value_type: ValueType::String,
341 required: false,
342 description: None,
343 constraints: None,
344 },
345 ),
346 (
347 "port".into(),
348 PropertyDef {
349 value_type: ValueType::Int,
350 required: false,
351 description: None,
352 constraints: None,
353 },
354 ),
355 ]),
356 subtypes: None,
357 parent_type: None,
358 },
359 ),
360 (
361 "signal".into(),
362 NodeTypeDef {
363 description: None,
364 properties: BTreeMap::new(),
365 subtypes: None,
366 parent_type: None,
367 },
368 ),
369 ]),
370 edge_types: BTreeMap::from([(
371 "RUNS_ON".into(),
372 EdgeTypeDef {
373 description: None,
374 source_types: vec!["entity".into()],
375 target_types: vec!["entity".into()],
376 properties: BTreeMap::new(),
377 },
378 )]),
379 }
380 }
381
382 fn sample_op() -> GraphOp {
383 GraphOp::AddNode {
384 node_id: "server-1".into(),
385 node_type: "entity".into(),
386 label: "Production Server".into(),
387 properties: BTreeMap::from([
388 ("ip".into(), Value::String("10.0.0.1".into())),
389 ("port".into(), Value::Int(8080)),
390 ]),
391 subtype: None,
392 }
393 }
394
395 fn sample_clock() -> LamportClock {
396 LamportClock::with_values("inst-a", 1, 0)
397 }
398
399 fn sample_entry() -> Entry {
400 Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a")
401 }
402
403 #[test]
404 fn entry_hash_deterministic() {
405 let e1 = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
406 let e2 = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
407 assert_eq!(e1.hash, e2.hash);
408 }
409
410 #[test]
411 fn entry_hash_changes_on_mutation() {
412 let e1 = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
413 let different_op = GraphOp::AddNode {
414 node_id: "server-2".into(),
415 node_type: "entity".into(),
416 label: "Other Server".into(),
417 properties: BTreeMap::new(),
418 subtype: None,
419 };
420 let e2 = Entry::new(different_op, vec![], vec![], sample_clock(), "inst-a");
421 assert_ne!(e1.hash, e2.hash);
422 }
423
424 #[test]
425 fn entry_hash_changes_with_different_author() {
426 let e1 = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
427 let e2 = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-b");
428 assert_ne!(e1.hash, e2.hash);
429 }
430
431 #[test]
432 fn entry_hash_changes_with_different_clock() {
433 let e1 = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
434 let mut clock2 = sample_clock();
435 clock2.physical_ms = 99;
436 let e2 = Entry::new(sample_op(), vec![], vec![], clock2, "inst-a");
437 assert_ne!(e1.hash, e2.hash);
438 }
439
440 #[test]
441 fn entry_hash_changes_with_different_next() {
442 let e1 = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
443 let e2 = Entry::new(
444 sample_op(),
445 vec![[0u8; 32]],
446 vec![],
447 sample_clock(),
448 "inst-a",
449 );
450 assert_ne!(e1.hash, e2.hash);
451 }
452
453 #[test]
454 fn entry_verify_hash_valid() {
455 let entry = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
456 assert!(entry.verify_hash());
457 }
458
459 #[test]
460 fn entry_verify_hash_reject_tampered() {
461 let mut entry = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
462 entry.author = "evil-node".into();
463 assert!(!entry.verify_hash());
464 }
465
466 #[test]
467 fn entry_roundtrip_msgpack() {
468 let entry = Entry::new(
469 sample_op(),
470 vec![[1u8; 32]],
471 vec![[2u8; 32]],
472 sample_clock(),
473 "inst-a",
474 );
475 let bytes = entry.to_bytes();
476 let decoded = Entry::from_bytes(&bytes).unwrap();
477 assert_eq!(entry, decoded);
478 }
479
480 #[test]
481 fn entry_next_links_causal() {
482 let e1 = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
483 let e2 = Entry::new(
484 GraphOp::RemoveNode {
485 node_id: "server-1".into(),
486 },
487 vec![e1.hash],
488 vec![],
489 LamportClock::with_values("inst-a", 2, 0),
490 "inst-a",
491 );
492 assert_eq!(e2.next, vec![e1.hash]);
493 assert!(e2.verify_hash());
494 }
495
496 #[test]
497 fn graphop_all_variants_serialize() {
498 let ops = vec![
499 GraphOp::DefineOntology {
500 ontology: sample_ontology(),
501 },
502 sample_op(),
503 GraphOp::AddEdge {
504 edge_id: "e1".into(),
505 edge_type: "RUNS_ON".into(),
506 source_id: "svc-1".into(),
507 target_id: "server-1".into(),
508 properties: BTreeMap::new(),
509 },
510 GraphOp::UpdateProperty {
511 entity_id: "server-1".into(),
512 key: "cpu".into(),
513 value: Value::Float(85.5),
514 },
515 GraphOp::RemoveNode {
516 node_id: "server-1".into(),
517 },
518 GraphOp::RemoveEdge {
519 edge_id: "e1".into(),
520 },
521 GraphOp::ExtendOntology {
522 extension: crate::ontology::OntologyExtension {
523 node_types: BTreeMap::from([(
524 "metric".into(),
525 NodeTypeDef {
526 description: Some("A metric observation".into()),
527 properties: BTreeMap::new(),
528 subtypes: None,
529 parent_type: None,
530 },
531 )]),
532 edge_types: BTreeMap::new(),
533 node_type_updates: BTreeMap::new(),
534 },
535 },
536 GraphOp::Checkpoint {
537 ops: vec![
538 GraphOp::DefineOntology {
539 ontology: sample_ontology(),
540 },
541 GraphOp::AddNode {
542 node_id: "n1".into(),
543 node_type: "entity".into(),
544 subtype: None,
545 label: "Node 1".into(),
546 properties: BTreeMap::new(),
547 },
548 ],
549 op_clocks: vec![(1, 0), (2, 0)],
550 compacted_at_physical_ms: 1000,
551 compacted_at_logical: 5,
552 },
553 ];
554 for op in ops {
555 let entry = Entry::new(op, vec![], vec![], sample_clock(), "inst-a");
556 let bytes = entry.to_bytes();
557 let decoded = Entry::from_bytes(&bytes).unwrap();
558 assert_eq!(entry, decoded);
559 }
560 }
561
562 #[test]
563 fn genesis_entry_contains_ontology() {
564 let ont = sample_ontology();
565 let genesis = Entry::new(
566 GraphOp::DefineOntology {
567 ontology: ont.clone(),
568 },
569 vec![],
570 vec![],
571 LamportClock::new("inst-a"),
572 "inst-a",
573 );
574 match &genesis.payload {
575 GraphOp::DefineOntology { ontology } => assert_eq!(ontology, &ont),
576 _ => panic!("genesis should be DefineOntology"),
577 }
578 assert!(genesis.next.is_empty(), "genesis has no predecessors");
579 assert!(genesis.verify_hash());
580 }
581
582 #[test]
583 fn value_all_variants_roundtrip() {
584 let values = vec![
585 Value::Null,
586 Value::Bool(true),
587 Value::Int(42),
588 Value::Float(3.14),
589 Value::String("hello".into()),
590 Value::List(vec![Value::Int(1), Value::String("two".into())]),
591 Value::Map(BTreeMap::from([("key".into(), Value::Bool(false))])),
592 ];
593 for val in values {
594 let bytes = rmp_serde::to_vec(&val).unwrap();
595 let decoded: Value = rmp_serde::from_slice(&bytes).unwrap();
596 assert_eq!(val, decoded);
597 }
598 }
599
600 #[test]
601 fn hash_hex_format() {
602 let entry = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
603 let hex = entry.hash_hex();
604 assert_eq!(hex.len(), 64);
605 assert!(hex.chars().all(|c| c.is_ascii_hexdigit()));
606 }
607
608 #[test]
609 fn unsigned_entry_has_no_signature() {
610 let entry = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
611 assert!(!entry.is_signed());
612 assert!(entry.signature.is_none());
613 }
614
615 #[test]
616 fn unsigned_entry_roundtrip_preserves_none_signature() {
617 let entry = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
618 let bytes = entry.to_bytes();
619 let decoded = Entry::from_bytes(&bytes).unwrap();
620 assert_eq!(decoded.signature, None);
621 assert!(decoded.verify_hash());
622 }
623
624 #[cfg(feature = "signing")]
625 mod signing_tests {
626 use super::*;
627
628 fn test_keypair() -> ed25519_dalek::SigningKey {
629 use rand::rngs::OsRng;
630 ed25519_dalek::SigningKey::generate(&mut OsRng)
631 }
632
633 #[test]
634 fn signed_entry_roundtrip() {
635 let key = test_keypair();
636 let entry =
637 Entry::new_signed(sample_op(), vec![], vec![], sample_clock(), "inst-a", &key);
638
639 assert!(entry.is_signed());
640 assert!(entry.verify_hash());
641
642 let public = key.verifying_key();
643 assert!(entry.verify_signature(&public));
644 }
645
646 #[test]
647 fn signed_entry_serialization_roundtrip() {
648 let key = test_keypair();
649 let entry =
650 Entry::new_signed(sample_op(), vec![], vec![], sample_clock(), "inst-a", &key);
651
652 let bytes = entry.to_bytes();
653 let decoded = Entry::from_bytes(&bytes).unwrap();
654
655 assert!(decoded.is_signed());
656 assert!(decoded.verify_hash());
657 assert!(decoded.verify_signature(&key.verifying_key()));
658 }
659
660 #[test]
661 fn wrong_key_fails_verification() {
662 let key1 = test_keypair();
663 let key2 = test_keypair();
664
665 let entry =
666 Entry::new_signed(sample_op(), vec![], vec![], sample_clock(), "inst-a", &key1);
667
668 assert!(entry.verify_signature(&key1.verifying_key()));
670 assert!(!entry.verify_signature(&key2.verifying_key()));
672 }
673
674 #[test]
675 fn tampered_hash_fails_both_checks() {
676 let key = test_keypair();
677 let mut entry =
678 Entry::new_signed(sample_op(), vec![], vec![], sample_clock(), "inst-a", &key);
679
680 entry.hash[0] ^= 0xFF;
682
683 assert!(!entry.verify_hash());
684 assert!(!entry.verify_signature(&key.verifying_key()));
685 }
686
687 #[test]
688 fn unsigned_entry_passes_signature_check() {
689 let key = test_keypair();
691 let entry = Entry::new(sample_op(), vec![], vec![], sample_clock(), "inst-a");
692
693 assert!(!entry.is_signed());
694 assert!(entry.verify_signature(&key.verifying_key())); }
696 }
697
698 #[test]
701 fn value_int_json_roundtrip_preserves_type() {
702 let val = Value::Int(1);
703 let json = serde_json::to_string(&val).unwrap();
704 let back: Value = serde_json::from_str(&json).unwrap();
705 assert_eq!(
706 back,
707 Value::Int(1),
708 "Int(1) -> JSON -> back should stay Int, got {:?}",
709 back
710 );
711 }
712
713 #[test]
714 fn value_float_json_roundtrip_preserves_type() {
715 let val = Value::Float(1.0);
716 let json = serde_json::to_string(&val).unwrap();
717 let back: Value = serde_json::from_str(&json).unwrap();
718 assert_eq!(
719 back,
720 Value::Float(1.0),
721 "Float(1.0) -> JSON -> back should stay Float, got {:?}",
722 back
723 );
724 }
725
726 #[test]
727 fn value_float_json_includes_decimal() {
728 let json = serde_json::to_string(&Value::Float(1.0)).unwrap();
731 assert!(
732 json.contains('.'),
733 "Float(1.0) must serialize with decimal point, got: {}",
734 json
735 );
736 }
737
738 #[test]
739 fn graphop_with_mixed_values_json_roundtrip() {
740 let mut props = BTreeMap::new();
741 props.insert("count".into(), Value::Int(42));
742 props.insert("ratio".into(), Value::Float(1.0));
743 props.insert("name".into(), Value::String("test".into()));
744
745 let op = GraphOp::UpdateProperty {
746 entity_id: "e1".into(),
747 key: "data".into(),
748 value: Value::Map(props),
749 };
750
751 let json = serde_json::to_string(&op).unwrap();
752 let back: GraphOp = serde_json::from_str(&json).unwrap();
753
754 let entry1 = Entry::new(op, vec![], vec![], sample_clock(), "a");
756 let entry2 = Entry::new(back, vec![], vec![], sample_clock(), "a");
757 assert_eq!(
758 entry1.hash, entry2.hash,
759 "JSON round-trip changed the hash!"
760 );
761 }
762
763 #[test]
769 fn legacy_eight_element_entry_still_deserializes() {
770 let entry = sample_entry();
771 let legacy = (
775 entry.hash,
776 entry.payload.clone(),
777 entry.next.clone(),
778 entry.refs.clone(),
779 entry.clock.clone(),
780 entry.author.clone(),
781 entry.signature.clone(),
782 None::<Hash>,
783 );
784 let bytes = rmp_serde::to_vec(&legacy).unwrap();
785 assert_eq!(bytes[0], 0x98, "legacy fixture is not 8 elements");
786
787 let restored = Entry::from_bytes(&bytes).expect("legacy entry must load");
788 assert_eq!(restored, entry);
789 assert!(restored.verify_hash());
790 }
791
792 #[test]
795 fn legacy_entries_in_a_sequence_stay_aligned() {
796 let entry = sample_entry();
797 let legacy = |e: &Entry| {
798 (
799 e.hash,
800 e.payload.clone(),
801 e.next.clone(),
802 e.refs.clone(),
803 e.clock.clone(),
804 e.author.clone(),
805 e.signature.clone(),
806 None::<Hash>,
807 )
808 };
809 let bytes = rmp_serde::to_vec(&vec![legacy(&entry), legacy(&entry)]).unwrap();
810 let restored: Vec<Entry> = rmp_serde::from_slice(&bytes).expect("legacy vec must load");
811 assert_eq!(restored.len(), 2);
812 assert!(restored.iter().all(|e| e.verify_hash()));
813 }
814
815 #[test]
823 fn entry_wire_format_is_a_positional_array_of_seven() {
824 let bytes = sample_entry().to_bytes();
825
826 assert_eq!(
828 bytes[0], 0x97,
829 "Entry is no longer a 7-element positional array. Changing the \
830 field count breaks every persisted store and every peer running \
831 an older build; add a deserialization shim for the old arity \
832 (as S4 did for the 8-element form) before changing this."
833 );
834 for name in [b"payload".as_slice(), b"signature".as_slice()] {
837 assert!(
838 !bytes.windows(name.len()).any(|w| w == name),
839 "expected positional encoding, found a field name on the wire"
840 );
841 }
842 }
843
844 #[test]
847 fn define_lens_roundtrips() {
848 let op = GraphOp::DefineLens {
849 transforms: vec![1, 2, 3, 4],
850 };
851 let entry = Entry::new(op.clone(), vec![], vec![], sample_clock(), "author");
852 let bytes = entry.to_bytes();
853 let restored = Entry::from_bytes(&bytes).unwrap();
854 assert_eq!(restored.payload, op);
855 assert!(restored.verify_hash());
856 }
857}