1#[cfg(all(test, feature = "arbitrary"))]
215mod conformance;
216mod storage;
217use commonware_runtime::buffer::paged::CacheRef;
218use commonware_utils::Array;
219use std::num::NonZeroUsize;
220pub use storage::{Checkpoint, Cursor, Freezer};
221use thiserror::Error;
222
223pub enum Identifier<'a, K: Array> {
225 Cursor(Cursor),
226 Key(&'a K),
227}
228
229#[derive(Debug, Error)]
231pub enum Error {
232 #[error("runtime error: {0}")]
233 Runtime(#[from] commonware_runtime::Error),
234 #[error("journal error: {0}")]
235 Journal(#[from] crate::journal::Error),
236 #[error("codec error: {0}")]
237 Codec(#[from] commonware_codec::Error),
238 #[error("checkpoint does not match stored data")]
239 CheckpointMismatch,
240}
241
242#[derive(Clone)]
244pub struct Config<C> {
245 pub key_partition: String,
247
248 pub key_write_buffer: NonZeroUsize,
250
251 pub key_page_cache: CacheRef,
253
254 pub value_partition: String,
256
257 pub value_compression: Option<u8>,
259
260 pub value_write_buffer: NonZeroUsize,
262
263 pub value_target_size: u64,
265
266 pub table_partition: String,
268
269 pub table_initial_size: u32,
271
272 pub table_resize_frequency: u8,
275
276 pub table_resize_chunk_size: u32,
278
279 pub table_replay_buffer: NonZeroUsize,
281
282 pub codec_config: C,
284}
285
286#[cfg(test)]
287mod tests {
288 use super::*;
289 use commonware_codec::DecodeExt;
290 use commonware_formatting::hex;
291 use commonware_macros::{test_group, test_traced};
292 use commonware_runtime::{
293 Blob, Metrics as _, ReadOptions, Runner, Storage, Supervisor as _, WriteOptions,
294 deterministic,
295 };
296 use commonware_utils::{NZU16, NZUsize, sequence::FixedBytes};
297 use rand::{Rng, RngExt as _};
298 use std::num::NonZeroU16;
299
300 fn test_key(key: &str) -> FixedBytes<64> {
301 let mut buf = [0u8; 64];
302 let key = key.as_bytes();
303 assert!(key.len() <= buf.len());
304 buf[..key.len()].copy_from_slice(key);
305 FixedBytes::decode(buf.as_ref()).unwrap()
306 }
307
308 const DEFAULT_WRITE_BUFFER: usize = 1024;
309 const DEFAULT_VALUE_TARGET_SIZE: u64 = 10 * 1024 * 1024;
310 const DEFAULT_TABLE_INITIAL_SIZE: u32 = 256;
311 const DEFAULT_TABLE_RESIZE_FREQUENCY: u8 = 4;
312 const DEFAULT_TABLE_RESIZE_CHUNK_SIZE: u32 = 128; const DEFAULT_TABLE_REPLAY_BUFFER: usize = 64 * 1024; const PAGE_SIZE: NonZeroU16 = NZU16!(1024);
315 const PAGE_CACHE_SIZE: NonZeroUsize = NZUsize!(10);
316
317 fn test_put_get(compression: Option<u8>) {
318 let executor = deterministic::Runner::default();
320 executor.start(|context| async move {
321 let cfg = Config {
323 key_partition: "test-key-index".into(),
324 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
325 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
326 value_partition: "test-value-journal".into(),
327 value_compression: compression,
328 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
329 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
330 table_partition: "test-table".into(),
331 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
332 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
333 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
334 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
335 codec_config: (),
336 };
337 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
338 context.child("storage"),
339 cfg.clone(),
340 None,
341 )
342 .await
343 .expect("Failed to initialize freezer");
344
345 let key = test_key("testkey");
346 let data = 42;
347
348 let value = freezer
350 .get(Identifier::Key(&key))
351 .await
352 .expect("Failed to check key");
353 assert!(value.is_none());
354
355 let (freezer, _) = freezer
357 .put(key.clone(), data)
358 .await
359 .expect("Failed to put data");
360
361 let value = freezer
363 .get(Identifier::Key(&key))
364 .await
365 .expect("Failed to get data")
366 .expect("Data not found");
367 assert_eq!(value, data);
368
369 let buffer = context.encode();
371 assert!(buffer.contains("gets_total 2"), "{}", buffer);
372 assert!(buffer.contains("puts_total 1"), "{}", buffer);
373 assert!(buffer.contains("unnecessary_reads_total 0"), "{}", buffer);
374
375 freezer.sync().await.expect("Failed to sync data");
377 });
378 }
379
380 #[test_traced]
381 fn test_put_get_no_compression() {
382 test_put_get(None);
383 }
384
385 #[test_traced]
386 fn test_put_get_compression() {
387 test_put_get(Some(3));
388 }
389
390 #[test_traced]
391 fn test_has() {
392 let executor = deterministic::Runner::default();
394 executor.start(|context| async move {
395 let cfg = Config {
397 key_partition: "test-key-index".into(),
398 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
399 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
400 value_partition: "test-value-journal".into(),
401 value_compression: None,
402 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
403 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
404 table_partition: "test-table".into(),
405 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
406 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
407 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
408 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
409 codec_config: (),
410 };
411 let freezer =
412 Freezer::<_, FixedBytes<64>, i32>::init(context.child("storage"), cfg, None)
413 .await
414 .expect("Failed to initialize freezer");
415
416 let key = test_key("testkey");
418 assert!(!freezer.has(&key).await.expect("Failed to check key"));
419
420 let (freezer, _) = freezer
422 .put(key.clone(), 42)
423 .await
424 .expect("Failed to put data");
425 assert!(freezer.has(&key).await.expect("Failed to check key"));
426
427 assert!(
429 !freezer
430 .has(&test_key("otherkey"))
431 .await
432 .expect("Failed to check key")
433 );
434
435 let buffer = context.encode();
437 assert!(buffer.contains("has_total 3"), "{}", buffer);
438 assert!(buffer.contains("gets_total 0"), "{}", buffer);
439 });
440 }
441
442 #[test_traced]
443 fn test_multiple_keys() {
444 let executor = deterministic::Runner::default();
446 executor.start(|context| async move {
447 let cfg = Config {
449 key_partition: "test-key-index".into(),
450 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
451 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
452 value_partition: "test-value-journal".into(),
453 value_compression: None,
454 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
455 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
456 table_partition: "test-table".into(),
457 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
458 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
459 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
460 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
461 codec_config: (),
462 };
463 let mut freezer = Freezer::<_, FixedBytes<64>, i32>::init(
464 context.child("storage"),
465 cfg.clone(),
466 None,
467 )
468 .await
469 .expect("Failed to initialize freezer");
470
471 let keys = vec![
473 (test_key("key1"), 1),
474 (test_key("key2"), 2),
475 (test_key("key3"), 3),
476 (test_key("key4"), 4),
477 (test_key("key5"), 5),
478 ];
479
480 for (key, data) in &keys {
481 (freezer, _) = freezer
482 .put(key.clone(), *data)
483 .await
484 .expect("Failed to put data");
485 }
486
487 for (key, data) in &keys {
489 let retrieved = freezer
490 .get(Identifier::Key(key))
491 .await
492 .expect("Failed to get data")
493 .expect("Data not found");
494 assert_eq!(retrieved, *data);
495 }
496 });
497 }
498
499 #[test_traced]
500 fn test_collision_handling() {
501 let executor = deterministic::Runner::default();
503 executor.start(|context| async move {
504 let cfg = Config {
506 key_partition: "test-key-index".into(),
507 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
508 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
509 value_partition: "test-value-journal".into(),
510 value_compression: None,
511 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
512 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
513 table_partition: "test-table".into(),
514 table_initial_size: 4, table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
516 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
517 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
518 codec_config: (),
519 };
520 let mut freezer = Freezer::<_, FixedBytes<64>, i32>::init(
521 context.child("storage"),
522 cfg.clone(),
523 None,
524 )
525 .await
526 .expect("Failed to initialize freezer");
527
528 let keys = vec![
530 (test_key("key1"), 1),
531 (test_key("key2"), 2),
532 (test_key("key3"), 3),
533 (test_key("key4"), 4),
534 (test_key("key5"), 5),
535 (test_key("key6"), 6),
536 (test_key("key7"), 7),
537 (test_key("key8"), 8),
538 ];
539
540 for (key, data) in &keys {
541 (freezer, _) = freezer
542 .put(key.clone(), *data)
543 .await
544 .expect("Failed to put data");
545 }
546
547 let (freezer, _) = freezer.sync().await.expect("Failed to sync");
549
550 for (key, data) in &keys {
552 let retrieved = freezer
553 .get(Identifier::Key(key))
554 .await
555 .expect("Failed to get data")
556 .expect("Data not found");
557 assert_eq!(retrieved, *data);
558 }
559
560 let buffer = context.encode();
562 assert!(buffer.contains("gets_total 8"), "{}", buffer);
563 assert!(buffer.contains("unnecessary_reads_total 5"), "{}", buffer);
564 });
565 }
566
567 #[test_traced]
568 fn test_restart() {
569 let executor = deterministic::Runner::default();
571 executor.start(|context| async move {
572 let cfg = Config {
573 key_partition: "test-key-index".into(),
574 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
575 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
576 value_partition: "test-value-journal".into(),
577 value_compression: None,
578 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
579 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
580 table_partition: "test-table".into(),
581 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
582 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
583 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
584 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
585 codec_config: (),
586 };
587
588 let checkpoint = {
590 let mut freezer = Freezer::<_, FixedBytes<64>, i32>::init(
591 context.child("first"),
592 cfg.clone(),
593 None,
594 )
595 .await
596 .expect("Failed to initialize freezer");
597
598 let keys = vec![
599 (test_key("persist1"), 100),
600 (test_key("persist2"), 200),
601 (test_key("persist3"), 300),
602 ];
603
604 for (key, data) in &keys {
605 (freezer, _) = freezer
606 .put(key.clone(), *data)
607 .await
608 .expect("Failed to put data");
609 }
610
611 freezer.close().await.expect("Failed to close freezer")
612 };
613
614 {
616 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
617 context.child("second"),
618 cfg.clone(),
619 Some(checkpoint),
620 )
621 .await
622 .expect("Failed to initialize freezer");
623
624 let keys = vec![
625 (test_key("persist1"), 100),
626 (test_key("persist2"), 200),
627 (test_key("persist3"), 300),
628 ];
629
630 for (key, data) in &keys {
631 let retrieved = freezer
632 .get(Identifier::Key(key))
633 .await
634 .expect("Failed to get data")
635 .expect("Data not found");
636 assert_eq!(retrieved, *data);
637 }
638 }
639 });
640 }
641
642 #[test_traced]
643 fn test_crash_consistency() {
644 let executor = deterministic::Runner::default();
646 executor.start(|context| async move {
647 let cfg = Config {
648 key_partition: "test-key-index".into(),
649 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
650 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
651 value_partition: "test-value-journal".into(),
652 value_compression: None,
653 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
654 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
655 table_partition: "test-table".into(),
656 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
657 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
658 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
659 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
660 codec_config: (),
661 };
662
663 let checkpoint = {
665 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
666 context.child("first"),
667 cfg.clone(),
668 None,
669 )
670 .await
671 .expect("Failed to initialize freezer");
672
673 let (freezer, _) = freezer
674 .put(test_key("committed1"), 1)
675 .await
676 .expect("Failed to put data");
677 let (freezer, _) = freezer
678 .put(test_key("committed2"), 2)
679 .await
680 .expect("Failed to put data");
681
682 let (freezer, _) = freezer.sync().await.expect("Failed to sync");
684
685 let (freezer, _) = freezer
687 .put(test_key("uncommitted1"), 3)
688 .await
689 .expect("Failed to put data");
690 let (freezer, _) = freezer
691 .put(test_key("uncommitted2"), 4)
692 .await
693 .expect("Failed to put data");
694
695 freezer.close().await.expect("Failed to close")
697 };
698
699 {
701 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
702 context.child("second"),
703 cfg.clone(),
704 Some(checkpoint),
705 )
706 .await
707 .expect("Failed to initialize freezer");
708
709 assert_eq!(
711 freezer
712 .get(Identifier::Key(&test_key("committed1")))
713 .await
714 .unwrap(),
715 Some(1)
716 );
717 assert_eq!(
718 freezer
719 .get(Identifier::Key(&test_key("committed2")))
720 .await
721 .unwrap(),
722 Some(2)
723 );
724
725 if let Some(val) = freezer
728 .get(Identifier::Key(&test_key("uncommitted1")))
729 .await
730 .unwrap()
731 {
732 assert_eq!(val, 3);
733 }
734 if let Some(val) = freezer
735 .get(Identifier::Key(&test_key("uncommitted2")))
736 .await
737 .unwrap()
738 {
739 assert_eq!(val, 4);
740 }
741 }
742 });
743 }
744
745 #[test_traced]
746 fn test_destroy() {
747 let executor = deterministic::Runner::default();
749 executor.start(|context| async move {
750 let cfg = Config {
752 key_partition: "test-key-index".into(),
753 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
754 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
755 value_partition: "test-value-journal".into(),
756 value_compression: None,
757 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
758 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
759 table_partition: "test-table".into(),
760 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
761 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
762 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
763 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
764 codec_config: (),
765 };
766 {
767 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
768 context.child("first"),
769 cfg.clone(),
770 None,
771 )
772 .await
773 .expect("Failed to initialize freezer");
774
775 let (freezer, _) = freezer
776 .put(test_key("destroy1"), 1)
777 .await
778 .expect("Failed to put data");
779 let (freezer, _) = freezer
780 .put(test_key("destroy2"), 2)
781 .await
782 .expect("Failed to put data");
783
784 freezer.destroy().await.expect("Failed to destroy freezer");
786 }
787
788 {
790 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
791 context.child("second"),
792 cfg.clone(),
793 None,
794 )
795 .await
796 .expect("Failed to initialize freezer");
797
798 assert!(
800 freezer
801 .get(Identifier::Key(&test_key("destroy1")))
802 .await
803 .unwrap()
804 .is_none()
805 );
806 assert!(
807 freezer
808 .get(Identifier::Key(&test_key("destroy2")))
809 .await
810 .unwrap()
811 .is_none()
812 );
813 }
814 });
815 }
816
817 #[test_traced]
818 fn test_partial_table_entry_write() {
819 let executor = deterministic::Runner::default();
821 executor.start(|context| async move {
822 let cfg = Config {
824 key_partition: "test-key-index".into(),
825 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
826 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
827 value_partition: "test-value-journal".into(),
828 value_compression: None,
829 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
830 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
831 table_partition: "test-table".into(),
832 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
833 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
834 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
835 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
836 codec_config: (),
837 };
838 let checkpoint = {
839 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
840 context.child("first"),
841 cfg.clone(),
842 None,
843 )
844 .await
845 .expect("Failed to initialize freezer");
846
847 let (freezer, _) = freezer.put(test_key("key1"), 42).await.unwrap();
848 let (freezer, _) = freezer.sync().await.unwrap();
849 freezer.close().await.unwrap()
850 };
851
852 {
854 let (blob, _) = context.open(&cfg.table_partition, b"table").await.unwrap();
855 blob.write_at(0, vec![0xFF; 10], WriteOptions::SYNC)
857 .await
858 .unwrap();
859 }
860
861 {
863 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
864 context.child("second"),
865 cfg.clone(),
866 Some(checkpoint),
867 )
868 .await
869 .expect("Failed to initialize freezer");
870
871 let result = freezer
874 .get(Identifier::Key(&test_key("key1")))
875 .await
876 .unwrap();
877 assert!(result.is_none() || result == Some(42));
878 }
879 });
880 }
881
882 #[test_traced]
883 fn test_table_entry_invalid_crc() {
884 let executor = deterministic::Runner::default();
886 executor.start(|context| async move {
887 let cfg = Config {
888 key_partition: "test-key-index".into(),
889 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
890 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
891 value_partition: "test-value-journal".into(),
892 value_compression: None,
893 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
894 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
895 table_partition: "test-table".into(),
896 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
897 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
898 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
899 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
900 codec_config: (),
901 };
902
903 let checkpoint = {
905 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
906 context.child("first"),
907 cfg.clone(),
908 None,
909 )
910 .await
911 .expect("Failed to initialize freezer");
912
913 let (freezer, _) = freezer.put(test_key("key1"), 42).await.unwrap();
914 let (freezer, _) = freezer.sync().await.unwrap();
915 freezer.close().await.unwrap()
916 };
917
918 {
920 let (blob, _) = context.open(&cfg.table_partition, b"table").await.unwrap();
921 let entry_data = blob.read_at(0, 24, ReadOptions::default()).await.unwrap();
923 let mut corrupted = entry_data.coalesce();
924 corrupted.as_mut()[20] ^= 0xFF;
926 blob.write_at(0, corrupted, WriteOptions::SYNC)
927 .await
928 .unwrap();
929 }
930
931 {
933 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
934 context.child("second"),
935 cfg.clone(),
936 Some(checkpoint),
937 )
938 .await
939 .expect("Failed to initialize freezer");
940
941 let result = freezer
943 .get(Identifier::Key(&test_key("key1")))
944 .await
945 .unwrap();
946 assert!(result.is_none() || result == Some(42));
948 }
949 });
950 }
951
952 #[test_traced]
953 fn test_table_extra_bytes() {
954 let executor = deterministic::Runner::default();
956 executor.start(|context| async move {
957 let cfg = Config {
958 key_partition: "test-key-index".into(),
959 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
960 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
961 value_partition: "test-value-journal".into(),
962 value_compression: None,
963 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
964 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
965 table_partition: "test-table".into(),
966 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
967 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
968 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
969 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
970 codec_config: (),
971 };
972
973 let checkpoint = {
975 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
976 context.child("first"),
977 cfg.clone(),
978 None,
979 )
980 .await
981 .expect("Failed to initialize freezer");
982
983 let (freezer, _) = freezer.put(test_key("key1"), 42).await.unwrap();
984 let (freezer, _) = freezer.sync().await.unwrap();
985 freezer.close().await.unwrap()
986 };
987
988 {
990 let (blob, size) = context.open(&cfg.table_partition, b"table").await.unwrap();
991 blob.write_at(size, hex!("0xdeadbeef").to_vec(), WriteOptions::SYNC)
993 .await
994 .unwrap();
995 }
996
997 {
999 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
1000 context.child("second"),
1001 cfg.clone(),
1002 Some(checkpoint),
1003 )
1004 .await
1005 .expect("Failed to initialize freezer");
1006
1007 assert_eq!(
1009 freezer
1010 .get(Identifier::Key(&test_key("key1")))
1011 .await
1012 .unwrap(),
1013 Some(42)
1014 );
1015
1016 let (freezer, _) = freezer.put(test_key("key2"), 43).await.unwrap();
1018 assert_eq!(
1019 freezer
1020 .get(Identifier::Key(&test_key("key2")))
1021 .await
1022 .unwrap(),
1023 Some(43)
1024 );
1025 }
1026 });
1027 }
1028
1029 #[test_traced]
1030 fn test_indexing_across_resizes() {
1031 let executor = deterministic::Runner::default();
1033 executor.start(|context| async move {
1034 let cfg = Config {
1036 key_partition: "test-key-index".into(),
1037 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1038 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
1039 value_partition: "test-value-journal".into(),
1040 value_compression: None,
1041 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1042 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
1043 table_partition: "test-table".into(),
1044 table_initial_size: 2, table_resize_frequency: 2, table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
1047 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
1048 codec_config: (),
1049 };
1050 let mut freezer =
1051 Freezer::<_, FixedBytes<64>, i32>::init(context.child("first"), cfg.clone(), None)
1052 .await
1053 .expect("Failed to initialize freezer");
1054
1055 let mut keys = Vec::new();
1058 for i in 0..1000 {
1059 let key = test_key(&format!("key{i}"));
1060 keys.push((key.clone(), i));
1061
1062 (freezer, _) = freezer.put(key, i).await.expect("Failed to put data");
1064 (freezer, _) = freezer.sync().await.expect("Failed to sync");
1065 }
1066
1067 for (key, value) in &keys {
1069 let retrieved = freezer
1070 .get(Identifier::Key(key))
1071 .await
1072 .expect("Failed to get data")
1073 .expect("Data not found");
1074 assert_eq!(retrieved, *value, "Value mismatch for key after resizes");
1075 }
1076
1077 let buffer = context.encode();
1081 assert!(buffer.contains("first_resizes_total 8"), "{}", buffer);
1082
1083 let checkpoint = freezer.close().await.expect("Failed to close");
1085 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
1086 context.child("second"),
1087 cfg.clone(),
1088 Some(checkpoint),
1089 )
1090 .await
1091 .expect("Failed to reinitialize freezer");
1092
1093 for (key, value) in &keys {
1095 let retrieved = freezer
1096 .get(Identifier::Key(key))
1097 .await
1098 .expect("Failed to get data")
1099 .expect("Data not found");
1100 assert_eq!(retrieved, *value, "Value mismatch for key after restart");
1101 }
1102 });
1103 }
1104
1105 #[test_traced]
1106 fn test_insert_during_resize() {
1107 let executor = deterministic::Runner::default();
1108 executor.start(|context| async move {
1109 let cfg = Config {
1110 key_partition: "test-key-index".into(),
1111 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1112 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
1113 value_partition: "test-value-journal".into(),
1114 value_compression: None,
1115 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1116 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
1117 table_partition: "test-table".into(),
1118 table_initial_size: 2,
1119 table_resize_frequency: 1,
1120 table_resize_chunk_size: 1, table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
1122 codec_config: (),
1123 };
1124 let mut freezer = Freezer::<_, FixedBytes<64>, i32>::init(
1125 context.child("storage"),
1126 cfg.clone(),
1127 None,
1128 )
1129 .await
1130 .unwrap();
1131
1132 (freezer, _) = freezer.put(test_key("key0"), 0).await.unwrap();
1135 (freezer, _) = freezer.put(test_key("key2"), 1).await.unwrap();
1136 (freezer, _) = freezer.sync().await.unwrap(); assert!(freezer.resizing().is_some());
1140
1141 (freezer, _) = freezer.put(test_key("key6"), 2).await.unwrap();
1144 assert!(context.encode().contains("unnecessary_writes_total 1"));
1145 assert_eq!(freezer.resizable(), 3);
1146
1147 (freezer, _) = freezer.put(test_key("key3"), 3).await.unwrap();
1150 assert!(context.encode().contains("unnecessary_writes_total 1"));
1151 assert_eq!(freezer.resizable(), 3);
1152
1153 (freezer, _) = freezer.sync().await.unwrap();
1155 assert!(freezer.resizing().is_none());
1156 assert_eq!(freezer.resizable(), 2);
1157
1158 (freezer, _) = freezer.put(test_key("key4"), 4).await.unwrap();
1161 (freezer, _) = freezer.put(test_key("key7"), 5).await.unwrap();
1162 (freezer, _) = freezer.sync().await.unwrap();
1163
1164 assert!(freezer.resizing().is_some());
1166
1167 let keys = ["key0", "key2", "key6", "key3", "key4", "key7"];
1169 for (i, k) in keys.iter().enumerate() {
1170 assert_eq!(
1171 freezer.get(Identifier::Key(&test_key(k))).await.unwrap(),
1172 Some(i as i32)
1173 );
1174 }
1175
1176 while freezer.resizing().is_some() {
1178 (freezer, _) = freezer.sync().await.unwrap();
1179 }
1180
1181 assert_eq!(freezer.resizable(), 0);
1183 });
1184 }
1185
1186 #[test_traced]
1187 fn test_resize_after_startup() {
1188 let executor = deterministic::Runner::default();
1189 executor.start(|context| async move {
1190 let cfg = Config {
1191 key_partition: "test-key-index".into(),
1192 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1193 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
1194 value_partition: "test-value-journal".into(),
1195 value_compression: None,
1196 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1197 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
1198 table_partition: "test-table".into(),
1199 table_initial_size: 2,
1200 table_resize_frequency: 1,
1201 table_resize_chunk_size: 1, table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
1203 codec_config: (),
1204 };
1205
1206 let checkpoint = {
1208 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
1209 context.child("first"),
1210 cfg.clone(),
1211 None,
1212 )
1213 .await
1214 .unwrap();
1215
1216 let (freezer, _) = freezer.put(test_key("key0"), 0).await.unwrap();
1219 let (freezer, _) = freezer.put(test_key("key2"), 1).await.unwrap();
1220 let (freezer, checkpoint) = freezer.sync().await.unwrap();
1221
1222 assert!(freezer.resizing().is_some());
1224
1225 checkpoint
1226 };
1227
1228 let mut freezer = Freezer::<_, FixedBytes<64>, i32>::init(
1230 context.child("second"),
1231 cfg.clone(),
1232 Some(checkpoint),
1233 )
1234 .await
1235 .unwrap();
1236 assert_eq!(freezer.resizable(), 1);
1237 assert_eq!(freezer.resizing(), None);
1238
1239 (freezer, _) = freezer.sync().await.unwrap();
1241 assert_eq!(freezer.resizing(), Some(1));
1242
1243 while freezer.resizing().is_some() {
1245 (freezer, _) = freezer.sync().await.unwrap();
1246 }
1247
1248 assert_eq!(freezer.resizable(), 0);
1250 });
1251 }
1252
1253 fn test_operations_and_restart(num_keys: usize) -> String {
1254 let executor = deterministic::Runner::default();
1255 executor.start(|mut context| async move {
1256 let cfg = Config {
1257 key_partition: "test-key-index".into(),
1258 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1259 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
1260 value_partition: "test-value-journal".into(),
1261 value_compression: None,
1262 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1263 value_target_size: 128, table_partition: "test-table".into(),
1265 table_initial_size: 8, table_resize_frequency: 2, table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
1268 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
1269 codec_config: (),
1270 };
1271 let mut freezer = Freezer::<_, FixedBytes<96>, FixedBytes<256>>::init(
1272 context.child("init").with_attribute("index", 1),
1273 cfg.clone(),
1274 None,
1275 )
1276 .await
1277 .expect("Failed to initialize freezer");
1278
1279 let mut pairs = Vec::new();
1281
1282 for _ in 0..num_keys {
1283 let mut key = [0u8; 96];
1285 context.fill_bytes(&mut key);
1286 let key = FixedBytes::<96>::new(key);
1287
1288 let mut value = [0u8; 256];
1290 context.fill_bytes(&mut value);
1291 let value = FixedBytes::<256>::new(value);
1292
1293 (freezer, _) = freezer
1295 .put(key.clone(), value.clone())
1296 .await
1297 .expect("Failed to put data");
1298 pairs.push((key, value));
1299
1300 if context.random_bool(0.1) {
1302 (freezer, _) = freezer.sync().await.expect("Failed to sync");
1303 }
1304 }
1305
1306 (freezer, _) = freezer.sync().await.expect("Failed to sync");
1308
1309 for (key, value) in &pairs {
1311 let retrieved = freezer
1312 .get(Identifier::Key(key))
1313 .await
1314 .expect("Failed to get data")
1315 .expect("Data not found");
1316 assert_eq!(&retrieved, value);
1317 }
1318
1319 for (key, _) in &pairs {
1321 assert!(
1322 freezer
1323 .get(Identifier::Key(key))
1324 .await
1325 .expect("Failed to check key")
1326 .is_some()
1327 );
1328 }
1329
1330 for _ in 0..10 {
1332 let mut key = [0u8; 96];
1333 context.fill_bytes(&mut key);
1334 let key = FixedBytes::<96>::new(key);
1335 assert!(
1336 freezer
1337 .get(Identifier::Key(&key))
1338 .await
1339 .expect("Failed to check key")
1340 .is_none()
1341 );
1342 }
1343
1344 let checkpoint = freezer.close().await.expect("Failed to close freezer");
1346
1347 let mut freezer = Freezer::<_, FixedBytes<96>, FixedBytes<256>>::init(
1349 context.child("init").with_attribute("index", 2),
1350 cfg.clone(),
1351 Some(checkpoint),
1352 )
1353 .await
1354 .expect("Failed to initialize freezer");
1355
1356 for (key, value) in &pairs {
1358 let retrieved = freezer
1359 .get(Identifier::Key(key))
1360 .await
1361 .expect("Failed to get data")
1362 .expect("Data not found");
1363 assert_eq!(&retrieved, value);
1364 }
1365
1366 for _ in 0..20 {
1368 let mut key = [0u8; 96];
1369 context.fill_bytes(&mut key);
1370 let key = FixedBytes::<96>::new(key);
1371
1372 let mut value = [0u8; 256];
1373 context.fill_bytes(&mut value);
1374 let value = FixedBytes::<256>::new(value);
1375
1376 (freezer, _) = freezer.put(key, value).await.expect("Failed to put data");
1377 }
1378
1379 for _ in 0..3 {
1381 (freezer, _) = freezer.sync().await.expect("Failed to sync");
1382
1383 for _ in 0..5 {
1385 let mut key = [0u8; 96];
1386 context.fill_bytes(&mut key);
1387 let key = FixedBytes::<96>::new(key);
1388
1389 let mut value = [0u8; 256];
1390 context.fill_bytes(&mut value);
1391 let value = FixedBytes::<256>::new(value);
1392
1393 (freezer, _) = freezer.put(key, value).await.expect("Failed to put data");
1394 }
1395 }
1396
1397 freezer.sync().await.expect("Failed to sync");
1399
1400 context.auditor().state()
1402 })
1403 }
1404
1405 #[test_group("slow")]
1406 #[test_traced]
1407 fn test_determinism() {
1408 let state1 = test_operations_and_restart(1_000);
1409 let state2 = test_operations_and_restart(1_000);
1410 assert_eq!(state1, state2);
1411 }
1412
1413 #[test_traced]
1414 fn test_put_multiple_updates() {
1415 let executor = deterministic::Runner::default();
1417 executor.start(|context| async move {
1418 let cfg = Config {
1420 key_partition: "test-key-index".into(),
1421 key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1422 key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
1423 value_partition: "test-value-journal".into(),
1424 value_compression: None,
1425 value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
1426 value_target_size: DEFAULT_VALUE_TARGET_SIZE,
1427 table_partition: "test-table".into(),
1428 table_initial_size: DEFAULT_TABLE_INITIAL_SIZE,
1429 table_resize_frequency: DEFAULT_TABLE_RESIZE_FREQUENCY,
1430 table_resize_chunk_size: DEFAULT_TABLE_RESIZE_CHUNK_SIZE,
1431 table_replay_buffer: NZUsize!(DEFAULT_TABLE_REPLAY_BUFFER),
1432 codec_config: (),
1433 };
1434 let freezer = Freezer::<_, FixedBytes<64>, i32>::init(
1435 context.child("storage"),
1436 cfg.clone(),
1437 None,
1438 )
1439 .await
1440 .expect("Failed to initialize freezer");
1441
1442 let key = test_key("key1");
1443
1444 let (freezer, _) = freezer
1445 .put(key.clone(), 1)
1446 .await
1447 .expect("Failed to put data");
1448 let (freezer, _) = freezer
1449 .put(key.clone(), 2)
1450 .await
1451 .expect("Failed to put data");
1452 let (freezer, _) = freezer.sync().await.expect("Failed to sync");
1453 assert_eq!(
1454 freezer
1455 .get(Identifier::Key(&key))
1456 .await
1457 .expect("Failed to get data")
1458 .unwrap(),
1459 2
1460 );
1461 });
1462 }
1463}