1use crate::types::Value;
2use serde::{Deserialize, Serialize};
3
4#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
5pub enum WalRecord {
6 InsertNode {
7 label: String,
8 key: String,
9 props: Vec<(String, Value)>,
10 },
11 InsertEdge {
12 edge_type: String,
13 src_key: String,
14 dst_key: String,
15 },
16 SetProp {
17 key: String,
18 field: String,
19 value: Value,
20 },
21 CreateRule {
22 def_bytes: Vec<u8>,
23 },
24 DeleteRule {
25 name: String,
26 },
27 RemoveProp {
29 key: String,
30 field: String,
31 },
32 DeleteEdge {
33 edge_type: String,
34 src_key: String,
35 dst_key: String,
36 },
37 DeleteNode {
38 key: String,
39 },
40 Batch(Vec<WalRecord>),
45 RebuildRule {
48 name: String,
49 },
50 CreateView {
54 def_bytes: Vec<u8>,
55 },
56 DeleteView {
58 name: String,
59 },
60 EnableFulltext {
65 label: String,
66 field: String,
67 },
68 DisableFulltext {
71 label: String,
72 field: String,
73 },
74 InsertNodeId {
78 label: u32,
79 key: String,
80 props: Vec<(u32, Value)>,
81 },
82 SetPropId {
84 id: u32,
85 field: u32,
86 value: Value,
87 },
88 InsertEdgeId {
90 etype: u32,
91 src: u32,
92 dst: u32,
93 },
94 Intern {
102 id: u32,
103 text: String,
104 },
105 DerivedEdgeAdded {
119 rule: String,
120 edge_type: String,
121 src_key: String,
122 dst_key: String,
123 },
124 DerivedEdgeRetracted {
126 rule: String,
127 edge_type: String,
128 src_key: String,
129 dst_key: String,
130 },
131 RenameNode {
138 old_key: String,
139 new_key: String,
140 },
141 EnableIndex {
148 label: String,
149 field: String,
150 },
151 DisableIndex {
154 label: String,
155 field: String,
156 },
157}
158
159pub fn encode_record(rec: &WalRecord) -> Vec<u8> {
165 if let WalRecord::Batch(inner) = rec {
166 debug_assert!(
167 !inner.iter().any(|r| matches!(r, WalRecord::Batch(_))),
168 "nested Batch is invalid: a Batch may not contain another Batch"
169 );
170 }
171 let payload = bincode::serialize(rec).expect("walrecord serialize cannot fail");
172 let crc = crc32fast::hash(&payload);
173 let mut out = Vec::with_capacity(8 + payload.len());
174 out.extend((payload.len() as u32).to_le_bytes());
175 out.extend(crc.to_le_bytes());
176 out.extend(payload);
177 out
178}
179
180pub fn wal_commits(bytes: &[u8]) -> u64 {
191 decode_all(bytes).0.len() as u64
192}
193
194pub fn decode_all(bytes: &[u8]) -> (Vec<WalRecord>, usize) {
203 let mut recs = Vec::new();
204 let mut pos = 0usize;
205 loop {
206 if bytes.len() < pos + 8 {
207 return (recs, pos);
208 }
209 let len = u32::from_le_bytes(bytes[pos..pos + 4].try_into().unwrap()) as usize;
212 let crc = u32::from_le_bytes(bytes[pos + 4..pos + 8].try_into().unwrap());
213 let start = pos + 8;
214 if bytes.len() < start + len {
215 return (recs, pos); }
217 let payload = &bytes[start..start + len];
218 if crc32fast::hash(payload) != crc {
219 return (recs, pos); }
221 match bincode::deserialize::<WalRecord>(payload) {
222 Ok(WalRecord::Batch(inner)) => {
223 if inner.iter().any(|r| matches!(r, WalRecord::Batch(_))) {
225 return (recs, pos);
226 }
227 recs.push(WalRecord::Batch(inner));
228 }
229 Ok(r) => recs.push(r),
230 Err(_) => return (recs, pos),
231 }
232 pos = start + len;
233 }
234}
235
236#[cfg(test)]
237mod tests {
238 use super::*;
239 use crate::types::Value;
240
241 fn sample() -> Vec<WalRecord> {
242 vec![
243 WalRecord::InsertNode {
244 label: "L".into(),
245 key: "k1".into(),
246 props: vec![("f".into(), Value::Int(1))],
247 },
248 WalRecord::InsertEdge {
249 edge_type: "E".into(),
250 src_key: "k1".into(),
251 dst_key: "k2".into(),
252 },
253 ]
254 }
255
256 #[test]
257 fn roundtrip_multiple_records() {
258 let mut bytes = Vec::new();
259 for r in sample() {
260 bytes.extend(encode_record(&r));
261 }
262 let (recs, consumed) = decode_all(&bytes);
263 assert_eq!(recs, sample());
264 assert_eq!(consumed, bytes.len());
265 }
266
267 #[test]
268 fn torn_tail_is_dropped_whole() {
269 let mut bytes = Vec::new();
270 for r in sample() {
271 bytes.extend(encode_record(&r));
272 }
273 let full = bytes.len();
274 let first = encode_record(&sample()[0]).len();
275 bytes.truncate(full - 3); let (recs, consumed) = decode_all(&bytes);
277 assert_eq!(recs.len(), 1);
278 assert_eq!(consumed, first);
279 }
280
281 #[test]
282 fn corrupt_crc_stops_replay_at_last_valid() {
283 let mut bytes = encode_record(&sample()[0]);
284 let n = bytes.len();
285 bytes[n - 1] ^= 0xFF; let (recs, consumed) = decode_all(&bytes);
287 assert!(recs.is_empty());
288 assert_eq!(consumed, 0);
289 }
290
291 #[test]
292 fn empty_input_is_fine() {
293 let (recs, consumed) = decode_all(&[]);
294 assert!(recs.is_empty());
295 assert_eq!(consumed, 0);
296 }
297
298 #[test]
301 fn roundtrip_remove_prop() {
302 let r = WalRecord::RemoveProp {
303 key: "n1".into(),
304 field: "age".into(),
305 };
306 let bytes = encode_record(&r);
307 let (recs, _) = decode_all(&bytes);
308 assert_eq!(recs, vec![r]);
309 }
310
311 #[test]
312 fn roundtrip_delete_edge() {
313 let r = WalRecord::DeleteEdge {
314 edge_type: "KNOWS".into(),
315 src_key: "a".into(),
316 dst_key: "b".into(),
317 };
318 let bytes = encode_record(&r);
319 let (recs, _) = decode_all(&bytes);
320 assert_eq!(recs, vec![r]);
321 }
322
323 #[test]
324 fn roundtrip_delete_node() {
325 let r = WalRecord::DeleteNode { key: "x".into() };
326 let bytes = encode_record(&r);
327 let (recs, _) = decode_all(&bytes);
328 assert_eq!(recs, vec![r]);
329 }
330
331 #[test]
332 fn roundtrip_rebuild_rule() {
333 let r = WalRecord::RebuildRule { name: "eq".into() };
334 let bytes = encode_record(&r);
335 let (recs, _) = decode_all(&bytes);
336 assert_eq!(recs, vec![r]);
337 }
338
339 #[test]
340 fn batch_of_three_is_one_frame() {
341 let inner = vec![
342 WalRecord::DeleteNode { key: "a".into() },
343 WalRecord::DeleteNode { key: "b".into() },
344 WalRecord::DeleteNode { key: "c".into() },
345 ];
346 let batch = WalRecord::Batch(inner.clone());
347 let frame = encode_record(&batch);
348
349 let (recs, consumed) = decode_all(&frame);
352 assert_eq!(consumed, frame.len(), "should consume the whole frame");
353 assert_eq!(recs.len(), 1, "one decoded record (the Batch)");
354 assert_eq!(recs[0], WalRecord::Batch(inner));
355 }
356
357 #[test]
358 fn torn_mid_batch_frame_drops_whole_batch() {
359 let pre = sample();
361 let batch = WalRecord::Batch(vec![
362 WalRecord::DeleteNode { key: "a".into() },
363 WalRecord::DeleteNode { key: "b".into() },
364 ]);
365 let mut bytes = Vec::new();
366 for r in &pre {
367 bytes.extend(encode_record(r));
368 }
369 let batch_start = bytes.len();
370 bytes.extend(encode_record(&batch));
371
372 bytes.truncate(bytes.len() - 3);
374
375 let (recs, consumed) = decode_all(&bytes);
376 assert_eq!(recs, pre, "only pre-batch records survive");
377 assert_eq!(
378 consumed, batch_start,
379 "valid_len stops at batch frame start"
380 );
381 }
382
383 #[test]
384 #[cfg(debug_assertions)]
385 #[should_panic(expected = "nested Batch")]
386 fn nested_batch_encode_panics_in_debug() {
387 let inner_batch = WalRecord::Batch(vec![WalRecord::DeleteNode { key: "z".into() }]);
388 let outer = WalRecord::Batch(vec![inner_batch]);
389 encode_record(&outer); }
391
392 #[test]
401 fn golden_bytes_pin_wire_format() {
402 let insert_node = WalRecord::InsertNode {
404 label: "L".into(),
405 key: "k".into(),
406 props: vec![],
407 };
408 #[rustfmt::skip]
409 let expected_insert_node: &[u8] = &[
410 30, 0, 0, 0, 114, 69, 253, 24,
412 0, 0, 0, 0,
414 1, 0, 0, 0, 0, 0, 0, 0, 76,
416 1, 0, 0, 0, 0, 0, 0, 0, 107,
418 0, 0, 0, 0, 0, 0, 0, 0,
420 ];
421 assert_eq!(
422 encode_record(&insert_node),
423 expected_insert_node,
424 "InsertNode wire format changed — this breaks all existing WAL files"
425 );
426
427 let remove_prop = WalRecord::RemoveProp {
430 key: "n1".into(),
431 field: "age".into(),
432 };
433 #[rustfmt::skip]
434 let expected_remove_prop: &[u8] = &[
435 25, 0, 0, 0, 35, 214, 55, 239,
437 5, 0, 0, 0,
439 2, 0, 0, 0, 0, 0, 0, 0, 110, 49,
441 3, 0, 0, 0, 0, 0, 0, 0, 97, 103, 101,
443 ];
444 assert_eq!(
445 encode_record(&remove_prop),
446 expected_remove_prop,
447 "RemoveProp wire format changed — this breaks all existing WAL files"
448 );
449
450 let batch_single = WalRecord::Batch(vec![WalRecord::DeleteNode { key: "z".into() }]);
458 #[rustfmt::skip]
461 let batch_payload: &[u8] = &[
462 8, 0, 0, 0,
464 1, 0, 0, 0, 0, 0, 0, 0,
466 7, 0, 0, 0,
468 1, 0, 0, 0, 0, 0, 0, 0, 122,
470 ];
471 let batch_crc = crc32fast::hash(batch_payload);
472 let mut expected_batch_frame: Vec<u8> = Vec::with_capacity(8 + batch_payload.len());
473 expected_batch_frame.extend((batch_payload.len() as u32).to_le_bytes());
474 expected_batch_frame.extend(batch_crc.to_le_bytes());
475 expected_batch_frame.extend_from_slice(batch_payload);
476 assert_eq!(
477 encode_record(&batch_single),
478 expected_batch_frame,
479 "Batch (discriminant 8) wire format changed — a variant may have \
480 been inserted before position 8, breaking all existing WAL Batch frames"
481 );
482
483 let rebuild = WalRecord::RebuildRule { name: "eq".into() };
486 #[rustfmt::skip]
487 let expected_rebuild: &[u8] = &[
488 14, 0, 0, 0, 242, 136, 144, 68,
490 9, 0, 0, 0,
492 2, 0, 0, 0, 0, 0, 0, 0, 101, 113,
494 ];
495 assert_eq!(
496 encode_record(&rebuild),
497 expected_rebuild,
498 "RebuildRule wire format changed — append-only WAL variants"
499 );
500 }
501
502 #[test]
505 fn roundtrip_create_view() {
506 let r = WalRecord::CreateView {
507 def_bytes: vec![1, 2, 3],
508 };
509 let bytes = encode_record(&r);
510 let (recs, _) = decode_all(&bytes);
511 assert_eq!(recs, vec![r]);
512 }
513
514 #[test]
515 fn roundtrip_delete_view() {
516 let r = WalRecord::DeleteView {
517 name: "my_view".into(),
518 };
519 let bytes = encode_record(&r);
520 let (recs, _) = decode_all(&bytes);
521 assert_eq!(recs, vec![r]);
522 }
523
524 #[test]
529 fn golden_bytes_pin_view_wire_format() {
530 let create_view = WalRecord::CreateView {
532 def_bytes: vec![0xDE, 0xAD],
533 };
534 let cv_payload = bincode::serialize(&create_view).unwrap();
535 assert_eq!(
537 &cv_payload[0..4],
538 &[10, 0, 0, 0],
539 "CreateView discriminant changed — a variant was inserted before position 10"
540 );
541
542 let delete_view = WalRecord::DeleteView { name: "v".into() };
544 let dv_payload = bincode::serialize(&delete_view).unwrap();
545 assert_eq!(
546 &dv_payload[0..4],
547 &[11, 0, 0, 0],
548 "DeleteView discriminant changed — a variant was inserted before position 11"
549 );
550
551 let mut buf = encode_record(&create_view);
553 buf.extend(encode_record(&delete_view));
554 let (recs, consumed) = decode_all(&buf);
555 assert_eq!(consumed, buf.len());
556 assert_eq!(recs.len(), 2);
557 assert_eq!(
558 recs[0],
559 WalRecord::CreateView {
560 def_bytes: vec![0xDE, 0xAD]
561 }
562 );
563 assert_eq!(recs[1], WalRecord::DeleteView { name: "v".into() });
564 }
565
566 #[test]
569 fn roundtrip_enable_fulltext() {
570 let r = WalRecord::EnableFulltext {
571 label: "Person".into(),
572 field: "bio".into(),
573 };
574 let bytes = encode_record(&r);
575 let (recs, _) = decode_all(&bytes);
576 assert_eq!(recs, vec![r]);
577 }
578
579 #[test]
580 fn roundtrip_disable_fulltext() {
581 let r = WalRecord::DisableFulltext {
582 label: "Person".into(),
583 field: "bio".into(),
584 };
585 let bytes = encode_record(&r);
586 let (recs, _) = decode_all(&bytes);
587 assert_eq!(recs, vec![r]);
588 }
589
590 #[test]
595 fn golden_bytes_pin_fulltext_wire_format() {
596 let enable = WalRecord::EnableFulltext {
598 label: "A".into(),
599 field: "b".into(),
600 };
601 let ep = bincode::serialize(&enable).unwrap();
602 assert_eq!(
603 &ep[0..4],
604 &[12, 0, 0, 0],
605 "EnableFulltext discriminant changed — a variant was inserted before position 12"
606 );
607
608 let disable = WalRecord::DisableFulltext {
610 label: "A".into(),
611 field: "b".into(),
612 };
613 let dp = bincode::serialize(&disable).unwrap();
614 assert_eq!(
615 &dp[0..4],
616 &[13, 0, 0, 0],
617 "DisableFulltext discriminant changed — a variant was inserted before position 13"
618 );
619
620 let mut buf = encode_record(&enable);
622 buf.extend(encode_record(&disable));
623 let (recs, consumed) = decode_all(&buf);
624 assert_eq!(consumed, buf.len());
625 assert_eq!(recs.len(), 2);
626 assert_eq!(
627 recs[0],
628 WalRecord::EnableFulltext {
629 label: "A".into(),
630 field: "b".into()
631 }
632 );
633 assert_eq!(
634 recs[1],
635 WalRecord::DisableFulltext {
636 label: "A".into(),
637 field: "b".into()
638 }
639 );
640 }
641
642 #[test]
647 fn golden_bytes_pin_property_index_wire_format() {
648 let enable = WalRecord::EnableIndex {
649 label: "A".into(),
650 field: "b".into(),
651 };
652 let ep = bincode::serialize(&enable).unwrap();
653 assert_eq!(
654 &ep[0..4],
655 &[21, 0, 0, 0],
656 "EnableIndex discriminant changed — a variant was inserted before position 21"
657 );
658
659 let disable = WalRecord::DisableIndex {
660 label: "A".into(),
661 field: "b".into(),
662 };
663 let dp = bincode::serialize(&disable).unwrap();
664 assert_eq!(
665 &dp[0..4],
666 &[22, 0, 0, 0],
667 "DisableIndex discriminant changed — a variant was inserted before position 22"
668 );
669
670 let mut buf = encode_record(&enable);
671 buf.extend(encode_record(&disable));
672 let (recs, consumed) = decode_all(&buf);
673 assert_eq!(consumed, buf.len());
674 assert_eq!(recs, vec![enable, disable]);
675 }
676
677 #[test]
678 fn roundtrip_dense_id_variants_append_after_fulltext() {
679 let recs = vec![
680 WalRecord::Intern {
681 id: 0,
682 text: "Person".into(),
683 },
684 WalRecord::InsertNodeId {
685 label: 0,
686 key: "a".into(),
687 props: vec![(1, Value::Int(1))],
688 },
689 WalRecord::SetPropId {
690 id: 0,
691 field: 1,
692 value: Value::Int(2),
693 },
694 WalRecord::InsertEdgeId {
695 etype: 2,
696 src: 0,
697 dst: 1,
698 },
699 ];
700 for r in &recs {
701 let bytes = encode_record(r);
702 let (got, n) = decode_all(&bytes);
703 assert_eq!(n, bytes.len());
704 assert_eq!(got, vec![r.clone()]);
705 }
706 let p = bincode::serialize(&recs[1]).unwrap();
707 assert_eq!(&p[0..4], &[14, 0, 0, 0], "InsertNodeId discriminant is 14");
708 let p = bincode::serialize(&recs[2]).unwrap();
709 assert_eq!(&p[0..4], &[15, 0, 0, 0], "SetPropId discriminant is 15");
710 let p = bincode::serialize(&recs[3]).unwrap();
711 assert_eq!(&p[0..4], &[16, 0, 0, 0], "InsertEdgeId discriminant is 16");
712 let p = bincode::serialize(&recs[0]).unwrap();
713 assert_eq!(&p[0..4], &[17, 0, 0, 0], "Intern discriminant is 17");
714 }
715
716 #[test]
717 fn history_marker_discriminants_pinned() {
718 let added = WalRecord::DerivedEdgeAdded {
722 rule: "r".into(),
723 edge_type: "T".into(),
724 src_key: "a".into(),
725 dst_key: "b".into(),
726 };
727 let retracted = WalRecord::DerivedEdgeRetracted {
728 rule: "r".into(),
729 edge_type: "T".into(),
730 src_key: "a".into(),
731 dst_key: "b".into(),
732 };
733 let pa = bincode::serialize(&added).unwrap();
734 assert_eq!(
735 &pa[0..4],
736 &[18, 0, 0, 0],
737 "DerivedEdgeAdded discriminant changed — a variant was inserted before position 18"
738 );
739 let pr = bincode::serialize(&retracted).unwrap();
740 assert_eq!(
741 &pr[0..4],
742 &[19, 0, 0, 0],
743 "DerivedEdgeRetracted discriminant changed — a variant was inserted before position 19"
744 );
745 let mut buf = encode_record(&added);
747 buf.extend(encode_record(&retracted));
748 let (recs, consumed) = decode_all(&buf);
749 assert_eq!(consumed, buf.len());
750 assert_eq!(recs.len(), 2);
751 assert_eq!(
752 recs[0],
753 WalRecord::DerivedEdgeAdded {
754 rule: "r".into(),
755 edge_type: "T".into(),
756 src_key: "a".into(),
757 dst_key: "b".into(),
758 }
759 );
760 assert_eq!(
761 recs[1],
762 WalRecord::DerivedEdgeRetracted {
763 rule: "r".into(),
764 edge_type: "T".into(),
765 src_key: "a".into(),
766 dst_key: "b".into(),
767 }
768 );
769 }
770
771 #[test]
776 fn rename_node_discriminant_pinned() {
777 let r = WalRecord::RenameNode {
778 old_key: "a".into(),
779 new_key: "b".into(),
780 };
781 let payload = bincode::serialize(&r).unwrap();
782 assert_eq!(
783 &payload[0..4],
784 &[20, 0, 0, 0],
785 "RenameNode discriminant changed — a variant was inserted before position 20"
786 );
787 let bytes = encode_record(&r);
789 let (recs, consumed) = decode_all(&bytes);
790 assert_eq!(consumed, bytes.len());
791 assert_eq!(recs.len(), 1);
792 assert_eq!(
793 recs[0],
794 WalRecord::RenameNode {
795 old_key: "a".into(),
796 new_key: "b".into(),
797 }
798 );
799 }
800
801 #[test]
802 fn nested_batch_decode_is_treated_as_corrupt() {
803 let inner_batch = WalRecord::Batch(vec![WalRecord::DeleteNode { key: "z".into() }]);
806 let outer = WalRecord::Batch(vec![inner_batch]);
807
808 let payload = bincode::serialize(&outer).unwrap();
810 let crc = crc32fast::hash(&payload);
811 let mut frame = Vec::with_capacity(8 + payload.len());
812 frame.extend((payload.len() as u32).to_le_bytes());
813 frame.extend(crc.to_le_bytes());
814 frame.extend(&payload);
815
816 let good = encode_record(&WalRecord::DeleteNode { key: "good".into() });
818 let good_len = good.len();
819 let mut bytes = good;
820 bytes.extend(&frame);
821
822 let (recs, consumed) = decode_all(&bytes);
823 assert_eq!(recs.len(), 1);
824 assert_eq!(recs[0], WalRecord::DeleteNode { key: "good".into() });
825 assert_eq!(
826 consumed, good_len,
827 "stops cleanly before the nested-batch frame"
828 );
829 }
830}