1#[cfg(all(test, feature = "arbitrary"))]
67mod conformance;
68mod storage;
69pub use storage::Metadata;
70use thiserror::Error;
71
72#[derive(Debug, Error)]
74pub enum Error {
75 #[error("runtime error: {0}")]
76 Runtime(#[from] commonware_runtime::Error),
77 #[error("corruption: {0}")]
78 Corruption(String),
79}
80
81#[derive(Clone)]
83pub struct Config<C> {
84 pub partition: String,
86
87 pub codec_config: C,
89}
90
91#[cfg(test)]
92mod tests {
93 use super::*;
94 use commonware_formatting::hex;
95 use commonware_macros::{test_group, test_traced};
96 use commonware_runtime::{
97 Blob, Metrics as _, ReadOptions, Runner, Storage, Supervisor as _, WriteOptions,
98 deterministic,
99 mocks::{
100 DelayedSyncContext, PendingSyncs, RecordingContext, Recordings, WriteFaultContext,
101 WriteFaults, drive_pending_syncs, fail_pending_syncs, release_pending_syncs,
102 },
103 };
104 use commonware_utils::sequence::U64;
105 use futures::FutureExt as _;
106 use rand::{Rng, RngExt as _};
107
108 fn assert_options(recordings: &Recordings, reads: &[ReadOptions], writes: &[WriteOptions]) {
109 let snapshot = recordings.snapshot();
110 assert_eq!(snapshot.reads.as_slice(), reads);
111 assert_eq!(snapshot.writes.as_slice(), writes);
112 recordings.clear();
113 }
114
115 fn assert_durability(pending: &PendingSyncs, calls: usize, starts: usize, completions: usize) {
116 assert_eq!(pending.calls(), calls);
117 assert_eq!(pending.starts(), starts);
118 assert_eq!(pending.completions(), completions);
119 }
120
121 #[test_traced]
122 fn test_io_options_and_durability() {
123 let executor = deterministic::Runner::default();
124 executor.start(|context| async move {
125 let pending = PendingSyncs::default();
126 let (recording, recordings) = RecordingContext::new(DelayedSyncContext {
127 inner: context,
128 pending: pending.clone(),
129 });
130 let cfg = Config {
131 partition: "test".into(),
132 codec_config: ((0..).into(), ()),
133 };
134 let key = U64::new(1);
135 let extra_key = U64::new(2);
136 let mut metadata =
137 Metadata::<_, U64, Vec<u8>>::init(recording.child("first"), cfg.clone())
138 .await
139 .unwrap();
140
141 metadata.put(key.clone(), vec![1; 8]);
143 metadata = metadata.sync().await.unwrap();
144 metadata = metadata.sync().await.unwrap();
145 recordings.clear();
146 pending.arm();
147
148 metadata.put(key.clone(), vec![2; 8]);
150 metadata = drive_pending_syncs(&pending, metadata.sync())
151 .await
152 .unwrap();
153 assert_options(
154 &recordings,
155 &[],
156 &[
157 WriteOptions::DONT_CACHE,
158 WriteOptions::DONT_CACHE,
159 WriteOptions::DONT_CACHE,
160 ],
161 );
162 assert_durability(&pending, 1, 0, 0);
163
164 metadata.put(key.clone(), vec![3; 8]);
166 let (next, handle) = metadata.start_sync().await.unwrap();
167 metadata = next;
168 drive_pending_syncs(&pending, handle).await.unwrap();
169 assert_options(
170 &recordings,
171 &[],
172 &[
173 WriteOptions::DONT_CACHE,
174 WriteOptions::DONT_CACHE,
175 WriteOptions::DONT_CACHE,
176 ],
177 );
178 assert_durability(&pending, 2, 1, 1);
179
180 metadata.put(extra_key.clone(), vec![4; 16]);
182 let (next, handle) = metadata.start_sync().await.unwrap();
183 metadata = next;
184 drive_pending_syncs(&pending, handle).await.unwrap();
185 assert_options(&recordings, &[], &[WriteOptions::DONT_CACHE]);
186 assert_durability(&pending, 3, 2, 2);
187
188 metadata = drive_pending_syncs(&pending, metadata.sync())
190 .await
191 .unwrap();
192 assert_options(
193 &recordings,
194 &[],
195 &[WriteOptions::SYNC | WriteOptions::DONT_CACHE],
196 );
197 assert_durability(&pending, 4, 2, 2);
198
199 metadata.remove(&extra_key);
201 let (next, handle) = metadata.start_sync().await.unwrap();
202 metadata = next;
203 drive_pending_syncs(&pending, handle).await.unwrap();
204 assert_options(&recordings, &[], &[WriteOptions::DONT_CACHE]);
205 assert_durability(&pending, 5, 3, 3);
206
207 metadata = drive_pending_syncs(&pending, metadata.sync())
209 .await
210 .unwrap();
211 assert_options(&recordings, &[], &[WriteOptions::DONT_CACHE]);
212 assert_durability(&pending, 6, 3, 3);
213
214 drop(metadata);
216 let metadata = Metadata::<_, U64, Vec<u8>>::init(recording.child("second"), cfg)
217 .await
218 .unwrap();
219 assert_options(
220 &recordings,
221 &[ReadOptions::DONT_CACHE, ReadOptions::DONT_CACHE],
222 &[],
223 );
224 metadata.destroy().await.unwrap();
225 });
226 }
227
228 #[test_traced]
229 fn test_start_sync_pipelined_destroy() {
230 let executor = deterministic::Runner::default();
231 executor.start(|context| async move {
232 let cfg = Config {
233 partition: "test".into(),
234 codec_config: ((0..).into(), ()),
235 };
236 let mut metadata =
237 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
238 .await
239 .unwrap();
240
241 let key = U64::new(1);
244 metadata.put(key.clone(), vec![3]);
245 let (mut metadata, h1) = metadata.start_sync().await.unwrap();
246 metadata.put(key.clone(), vec![4]);
247 let (metadata, h2) = metadata.start_sync().await.unwrap();
248 h1.await.unwrap();
249 h2.await.unwrap();
250 metadata.destroy().await.unwrap();
251
252 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
254 .await
255 .unwrap();
256 assert_eq!(metadata.get(&key), None, "destroyed store must be empty");
257 });
258 }
259
260 #[test_traced]
261 fn test_start_sync_failure_fails_next_sync() {
262 let executor = deterministic::Runner::default();
263 executor.start(|context| async move {
264 let pending = PendingSyncs::default();
265 let cfg = Config {
266 partition: "test".into(),
267 codec_config: ((0..).into(), ()),
268 };
269 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(
270 DelayedSyncContext {
271 inner: context.child("first"),
272 pending: pending.clone(),
273 },
274 cfg,
275 )
276 .await
277 .unwrap();
278
279 metadata.put(U64::new(1), vec![3]);
280 let (mut metadata, handle) = metadata.start_sync().await.unwrap();
281 fail_pending_syncs(&pending);
282 assert!(handle.await.is_err());
283
284 metadata.put(U64::new(1), vec![4]);
287 assert!(metadata.start_sync().await.is_err());
288 });
289 }
290
291 #[test_traced]
292 fn test_start_sync_dropped_handle_does_not_cancel() {
293 let executor = deterministic::Runner::default();
294 executor.start(|context| async move {
295 let pending = PendingSyncs::default();
296 let cfg = Config {
297 partition: "test".into(),
298 codec_config: ((0..).into(), ()),
299 };
300 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(
301 DelayedSyncContext {
302 inner: context.child("first"),
303 pending: pending.clone(),
304 },
305 cfg.clone(),
306 )
307 .await
308 .unwrap();
309
310 let key = U64::new(1);
312 metadata.put(key.clone(), vec![3]);
313 let (metadata, handle) = metadata.start_sync().await.unwrap();
314 drop(handle);
315 release_pending_syncs(&pending);
316 let metadata = drive_pending_syncs(&pending, metadata.sync())
317 .await
318 .unwrap();
319 drop(metadata);
320
321 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
322 .await
323 .unwrap();
324 assert_eq!(metadata.get(&key), Some(&vec![3]));
325 });
326 }
327
328 #[test_traced]
329 fn test_start_sync_dropped_handle_fails_next_sync() {
330 let executor = deterministic::Runner::default();
331 executor.start(|context| async move {
332 let pending = PendingSyncs::default();
333 let cfg = Config {
334 partition: "test".into(),
335 codec_config: ((0..).into(), ()),
336 };
337 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(
338 DelayedSyncContext {
339 inner: context.child("first"),
340 pending: pending.clone(),
341 },
342 cfg,
343 )
344 .await
345 .unwrap();
346
347 metadata.put(U64::new(1), vec![3]);
350 let (metadata, handle) = metadata.start_sync().await.unwrap();
351 drop(handle);
352 fail_pending_syncs(&pending);
353
354 assert!(metadata.sync().await.is_err());
355 });
356 }
357
358 #[test_traced]
359 fn test_start_sync_newest_copy_wins_on_reopen() {
360 let executor = deterministic::Runner::default();
361 executor.start(|context| async move {
362 let cfg = Config {
363 partition: "test".into(),
364 codec_config: ((0..).into(), ()),
365 };
366 let mut metadata =
367 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
368 .await
369 .unwrap();
370
371 let key = U64::new(1);
372 metadata.put(key.clone(), vec![3]);
373 let (mut metadata, h1) = metadata.start_sync().await.unwrap();
374 metadata.put(key.clone(), vec![4]);
375 let (metadata, h2) = metadata.start_sync().await.unwrap();
376 h1.await.unwrap();
377 h2.await.unwrap();
378 drop(metadata);
379
380 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
381 .await
382 .unwrap();
383 assert_eq!(metadata.get(&key), Some(&vec![4]));
384 });
385 }
386
387 #[test_traced]
388 fn test_start_sync_write_failure_consumes() {
389 let executor = deterministic::Runner::default();
390 executor.start(|context| async move {
391 let faults = WriteFaults::default();
392 let cfg = Config {
393 partition: "test".into(),
394 codec_config: ((0..).into(), ()),
395 };
396 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(
397 WriteFaultContext {
398 inner: context.child("first"),
399 faults: faults.clone(),
400 },
401 cfg.clone(),
402 )
403 .await
404 .unwrap();
405
406 let key = U64::new(1);
409 metadata.put(key.clone(), vec![1; 8]);
410 metadata = metadata.sync().await.unwrap();
411 metadata.put(key.clone(), vec![2; 8]);
412 metadata = metadata.sync().await.unwrap();
413
414 faults.arm();
417 metadata.put(key.clone(), vec![3; 8]);
418 assert!(metadata.start_sync().await.is_err());
419 faults.disarm();
420
421 let metadata = Metadata::<_, U64, Vec<u8>>::init(
424 WriteFaultContext {
425 inner: context.child("second"),
426 faults,
427 },
428 cfg,
429 )
430 .await
431 .unwrap();
432 assert_eq!(metadata.get(&key), Some(&vec![2; 8]));
433 });
434 }
435
436 #[test_traced]
437 fn test_start_sync_second_sync_waits_for_first() {
438 let executor = deterministic::Runner::default();
439 executor.start(|context| async move {
440 let pending = PendingSyncs::default();
441 let cfg = Config {
442 partition: "test".into(),
443 codec_config: ((0..).into(), ()),
444 };
445 let faults = WriteFaults::default();
446 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(
447 WriteFaultContext {
448 inner: DelayedSyncContext {
449 inner: context.child("first"),
450 pending: pending.clone(),
451 },
452 faults: faults.clone(),
453 },
454 cfg,
455 )
456 .await
457 .unwrap();
458
459 let key = U64::new(1);
463 metadata.put(key.clone(), vec![3]);
464 let (mut metadata, handle) = metadata.start_sync().await.unwrap();
465 metadata.put(key.clone(), vec![4]);
466 let writes_before = faults.writes();
467 let mut second = Box::pin(metadata.sync());
468 for _ in 0..8 {
469 assert!((&mut second).now_or_never().is_none());
470 }
471 assert_eq!(faults.writes(), writes_before, "second sync must not write");
472 assert_eq!(
473 pending.starts(),
474 1,
475 "second sync must not start a blob sync"
476 );
477
478 release_pending_syncs(&pending);
479 handle.await.unwrap();
480 let metadata = drive_pending_syncs(&pending, second).await.unwrap();
481 metadata.destroy().await.unwrap();
482 });
483 }
484
485 #[test_traced]
486 fn test_put_get_clear() {
487 let executor = deterministic::Runner::default();
489 executor.start(|context| async move {
490 let cfg = Config {
492 partition: "test".into(),
493 codec_config: ((0..).into(), ()),
494 };
495 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg)
496 .await
497 .unwrap();
498
499 let key = U64::new(42);
501 let value = metadata.get(&key);
502 assert!(value.is_none());
503
504 let buffer = context.encode();
506 assert!(buffer.contains("first_sync_rewrites_total 0"));
507 assert!(buffer.contains("first_sync_overwrites_total 0"));
508 assert!(buffer.contains("first_keys 0"));
509
510 let hello = b"hello".to_vec();
512 metadata.put(key.clone(), hello.clone());
513
514 let value = metadata.get(&key).unwrap();
516 assert_eq!(value, &hello);
517
518 let buffer = context.encode();
520 assert!(buffer.contains("first_sync_rewrites_total 0"));
521 assert!(buffer.contains("first_sync_overwrites_total 0"));
522 assert!(buffer.contains("first_keys 1"));
523
524 metadata = metadata.sync().await.unwrap();
526
527 let buffer = context.encode();
529 assert!(buffer.contains("first_sync_rewrites_total 1"));
530 assert!(buffer.contains("first_sync_overwrites_total 0"));
531 assert!(buffer.contains("first_keys 1"));
532
533 drop(metadata);
535 let cfg = Config {
536 partition: "test".into(),
537 codec_config: ((0..).into(), ()),
538 };
539 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
540 .await
541 .unwrap();
542
543 let buffer = context.encode();
545 assert!(buffer.contains("second_sync_rewrites_total 0"));
546 assert!(buffer.contains("second_sync_overwrites_total 0"));
547 assert!(buffer.contains("second_keys 1"));
548
549 let value = metadata.get(&key).unwrap();
551 assert_eq!(value, &hello);
552
553 metadata.clear();
555 let value = metadata.get(&key);
556 assert!(value.is_none());
557
558 let buffer = context.encode();
560 assert!(buffer.contains("second_sync_rewrites_total 0"));
561 assert!(buffer.contains("second_sync_overwrites_total 0"));
562 assert!(buffer.contains("second_keys 0"));
563
564 metadata.destroy().await.unwrap();
565 });
566 }
567
568 #[test_traced]
569 fn test_put_returns_previous_value() {
570 let executor = deterministic::Runner::default();
571 executor.start(|context| async move {
572 let cfg = Config {
573 partition: "test".into(),
574 codec_config: ((0..).into(), ()),
575 };
576 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg)
577 .await
578 .unwrap();
579
580 let key = U64::new(42);
581
582 let previous = metadata.put(key.clone(), b"first".to_vec());
584 assert!(previous.is_none());
585
586 let previous = metadata.put(key.clone(), b"second".to_vec());
588 assert_eq!(previous, Some(b"first".to_vec()));
589
590 let previous = metadata.put(key.clone(), b"third".to_vec());
592 assert_eq!(previous, Some(b"second".to_vec()));
593
594 assert_eq!(metadata.get(&key), Some(&b"third".to_vec()));
596
597 let other_key = U64::new(99);
599 let previous = metadata.put(other_key.clone(), b"other".to_vec());
600 assert!(previous.is_none());
601
602 metadata.sync().await.unwrap();
604
605 let cfg = Config {
606 partition: "test".into(),
607 codec_config: ((0..).into(), ()),
608 };
609 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
610 .await
611 .unwrap();
612
613 let previous = metadata.put(key.clone(), b"fourth".to_vec());
615 assert_eq!(previous, Some(b"third".to_vec()));
616
617 metadata.destroy().await.unwrap();
618 });
619 }
620
621 #[test_traced]
622 fn test_multi_sync() {
623 let executor = deterministic::Runner::default();
625 executor.start(|context| async move {
626 let cfg = Config {
628 partition: "test".into(),
629 codec_config: ((0..).into(), ()),
630 };
631 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg)
632 .await
633 .unwrap();
634
635 let key = U64::new(42);
637 let hello = b"hello".to_vec();
638 metadata.put(key.clone(), hello.clone());
639
640 metadata = metadata.sync().await.unwrap();
642
643 let buffer = context.encode();
645 assert!(buffer.contains("first_sync_rewrites_total 1"));
646 assert!(buffer.contains("first_sync_overwrites_total 0"));
647 assert!(buffer.contains("first_keys 1"));
648
649 let world = b"world".to_vec();
651 metadata.put(key.clone(), world.clone());
652 let key2 = U64::new(43);
653 let foo = b"foo".to_vec();
654 metadata.put(key2.clone(), foo.clone());
655
656 metadata = metadata.sync().await.unwrap();
658
659 let buffer = context.encode();
661 assert!(buffer.contains("first_sync_rewrites_total 2"));
662 assert!(buffer.contains("first_sync_overwrites_total 0"));
663 assert!(buffer.contains("first_keys 2"));
664
665 drop(metadata);
667 let cfg = Config {
668 partition: "test".into(),
669 codec_config: ((0..).into(), ()),
670 };
671 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
672 .await
673 .unwrap();
674
675 let buffer = context.encode();
677 assert!(buffer.contains("second_sync_rewrites_total 0"));
678 assert!(buffer.contains("second_sync_overwrites_total 0"));
679 assert!(buffer.contains("second_keys 2"));
680
681 let value = metadata.get(&key).unwrap();
683 assert_eq!(value, &world);
684 let value = metadata.get(&key2).unwrap();
685 assert_eq!(value, &foo);
686
687 metadata.remove(&key);
689
690 metadata = metadata.sync().await.unwrap();
692
693 let buffer = context.encode();
695 assert!(buffer.contains("second_sync_rewrites_total 1"));
696 assert!(buffer.contains("second_sync_overwrites_total 0"));
697 assert!(buffer.contains("second_keys 1"));
698
699 drop(metadata);
701 let cfg = Config {
702 partition: "test".into(),
703 codec_config: ((0..).into(), ()),
704 };
705 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("third"), cfg)
706 .await
707 .unwrap();
708
709 let buffer = context.encode();
711 assert!(buffer.contains("third_sync_rewrites_total 0"));
712 assert!(buffer.contains("third_sync_overwrites_total 0"));
713 assert!(buffer.contains("third_keys 1"));
714
715 let value = metadata.get(&key);
717 assert!(value.is_none());
718 let value = metadata.get(&key2).unwrap();
719 assert_eq!(value, &foo);
720
721 metadata.destroy().await.unwrap();
722 });
723 }
724
725 #[test_traced]
726 fn test_recover_corrupted_one() {
727 let executor = deterministic::Runner::default();
729 executor.start(|context| async move {
730 let cfg = Config {
732 partition: "test".into(),
733 codec_config: ((0..).into(), ()),
734 };
735 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg)
736 .await
737 .unwrap();
738
739 let key = U64::new(42);
741 let hello = b"hello".to_vec();
742 metadata.put(key.clone(), hello.clone());
743
744 metadata = metadata.sync().await.unwrap();
746
747 let world = b"world".to_vec();
749 metadata.put(key.clone(), world.clone());
750 let key2 = U64::new(43);
751 let foo = b"foo".to_vec();
752 metadata.put(key2, foo.clone());
753
754 metadata.sync().await.unwrap();
756
757 let (blob, _) = context.open("test", b"left").await.unwrap();
759 blob.write_at(0, b"corrupted".to_vec(), WriteOptions::SYNC)
760 .await
761 .unwrap();
762
763 let cfg = Config {
765 partition: "test".into(),
766 codec_config: ((0..).into(), ()),
767 };
768 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
769 .await
770 .unwrap();
771
772 let value = metadata.get(&key).unwrap();
774 assert_eq!(value, &hello);
775
776 metadata.destroy().await.unwrap();
777 });
778 }
779
780 #[test_traced]
781 fn test_recovered_mirror_supports_shrinking_rewrite() {
782 let executor = deterministic::Runner::default();
783 executor.start(|context| async move {
784 let cfg = Config {
785 partition: "test".into(),
786 codec_config: ((0..).into(), ()),
787 };
788 let mut metadata =
789 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
790 .await
791 .unwrap();
792 let key = U64::new(42);
793 let hello = b"hello".to_vec();
794 metadata.put(key.clone(), hello.clone());
795 metadata = metadata.sync().await.unwrap();
796 metadata.put(key.clone(), b"world".to_vec());
797 metadata.put(U64::new(43), b"foo".to_vec());
798 metadata.sync().await.unwrap();
799
800 let (blob, _) = context.open("test", b"left").await.unwrap();
802 blob.write_at(0, b"corrupted".to_vec(), WriteOptions::SYNC)
803 .await
804 .unwrap();
805
806 let mut metadata =
809 Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg.clone())
810 .await
811 .unwrap();
812 assert_eq!(metadata.get(&key).unwrap(), &hello);
813 metadata.clear();
814 metadata.sync().await.unwrap();
815
816 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("third"), cfg)
817 .await
818 .unwrap();
819 assert!(metadata.get(&key).is_none());
820 metadata.destroy().await.unwrap();
821 });
822 }
823
824 #[test_traced]
825 fn test_recover_corrupted_both() {
826 let executor = deterministic::Runner::default();
828 executor.start(|context| async move {
829 let cfg = Config {
831 partition: "test".into(),
832 codec_config: ((0..).into(), ()),
833 };
834 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg)
835 .await
836 .unwrap();
837
838 let key = U64::new(42);
840 let hello = b"hello".to_vec();
841 metadata.put(key.clone(), hello.clone());
842
843 metadata = metadata.sync().await.unwrap();
845
846 let world = b"world".to_vec();
848 metadata.put(key.clone(), world.clone());
849 let key2 = U64::new(43);
850 let foo = b"foo".to_vec();
851 metadata.put(key2, foo.clone());
852
853 metadata.sync().await.unwrap();
855
856 let (blob, _) = context.open("test", b"left").await.unwrap();
858 blob.write_at(0, b"corrupted".to_vec(), WriteOptions::SYNC)
859 .await
860 .unwrap();
861 let (blob, _) = context.open("test", b"right").await.unwrap();
862 blob.write_at(0, b"corrupted".to_vec(), WriteOptions::SYNC)
863 .await
864 .unwrap();
865
866 for child in ["second", "third"] {
869 let cfg = Config {
870 partition: "test".into(),
871 codec_config: ((0..).into(), ()),
872 };
873 let result = Metadata::<_, U64, Vec<u8>>::init(context.child(child), cfg).await;
874 assert!(matches!(result, Err(Error::Corruption(_))));
875 }
876 });
877 }
878
879 #[test_traced]
880 fn test_recover_corrupted_truncate() {
881 let executor = deterministic::Runner::default();
883 executor.start(|context| async move {
884 let cfg = Config {
886 partition: "test".into(),
887 codec_config: ((0..).into(), ()),
888 };
889 let mut metadata = Metadata::init(context.child("first"), cfg).await.unwrap();
890
891 let key = U64::new(42);
893 let hello = b"hello".to_vec();
894 metadata.put(key.clone(), hello.clone());
895
896 metadata = metadata.sync().await.unwrap();
898
899 let world = b"world".to_vec();
901 metadata.put(key.clone(), world.clone());
902 let key2 = U64::new(43);
903 let foo = b"foo".to_vec();
904 metadata.put(key2, foo.clone());
905
906 metadata.sync().await.unwrap();
908
909 let (blob, len) = context.open("test", b"left").await.unwrap();
911 blob.resize(len - 8).await.unwrap();
912 blob.sync().await.unwrap();
913
914 let cfg = Config {
916 partition: "test".into(),
917 codec_config: ((0..).into(), ()),
918 };
919 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
920 .await
921 .unwrap();
922
923 let value = metadata.get(&key).unwrap();
925 assert_eq!(value, &hello);
926
927 metadata.destroy().await.unwrap();
928 });
929 }
930
931 #[test_traced]
932 fn test_recover_corrupted_short() {
933 let executor = deterministic::Runner::default();
935 executor.start(|context| async move {
936 let cfg = Config {
938 partition: "test".into(),
939 codec_config: ((0..).into(), ()),
940 };
941 let mut metadata = Metadata::init(context.child("first"), cfg).await.unwrap();
942
943 let key = U64::new(42);
945 let hello = b"hello".to_vec();
946 metadata.put(key.clone(), hello.clone());
947
948 metadata = metadata.sync().await.unwrap();
950
951 let world = b"world".to_vec();
953 metadata.put(key.clone(), world.clone());
954 let key2 = U64::new(43);
955 let foo = b"foo".to_vec();
956 metadata.put(key2, foo.clone());
957
958 metadata.sync().await.unwrap();
960
961 let (blob, _) = context.open("test", b"left").await.unwrap();
963 blob.resize(5).await.unwrap();
964 blob.sync().await.unwrap();
965
966 let cfg = Config {
968 partition: "test".into(),
969 codec_config: ((0..).into(), ()),
970 };
971 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
972 .await
973 .unwrap();
974
975 let value = metadata.get(&key).unwrap();
977 assert_eq!(value, &hello);
978
979 metadata.destroy().await.unwrap();
980 });
981 }
982
983 #[test_traced]
984 fn test_unclean_shutdown() {
985 let executor = deterministic::Runner::default();
987 executor.start(|context| async move {
988 let key = U64::new(42);
989 let hello = b"hello".to_vec();
990 {
991 let cfg = Config {
993 partition: "test".into(),
994 codec_config: ((0..).into(), ()),
995 };
996 let mut metadata = Metadata::init(context.child("first"), cfg).await.unwrap();
997
998 metadata.put(key.clone(), hello.clone());
1000
1001 }
1003
1004 let cfg = Config {
1006 partition: "test".into(),
1007 codec_config: ((0..).into(), ()),
1008 };
1009 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
1010 .await
1011 .unwrap();
1012
1013 let value = metadata.get(&key);
1015 assert!(value.is_none());
1016
1017 let buffer = context.encode();
1019 assert!(buffer.contains("second_sync_rewrites_total 0"));
1020 assert!(buffer.contains("second_sync_overwrites_total 0"));
1021 assert!(buffer.contains("second_keys 0"));
1022
1023 metadata.destroy().await.unwrap();
1024 });
1025 }
1026
1027 #[test_traced]
1028 #[should_panic(expected = "usize value is larger than u32")]
1029 fn test_value_too_big_error() {
1030 let executor = deterministic::Runner::default();
1032 executor.start(|context| async move {
1033 let cfg = Config {
1035 partition: "test".into(),
1036 codec_config: ((0..).into(), ()),
1037 };
1038 let mut metadata = Metadata::init(context.child("storage"), cfg).await.unwrap();
1039
1040 let value = vec![0u8; (u32::MAX as usize) + 1];
1042 metadata.put(U64::new(1), value);
1043
1044 metadata.sync().await.unwrap();
1046 });
1047 }
1048
1049 #[test_traced]
1050 fn test_delta_writes() {
1051 let executor = deterministic::Runner::default();
1053 executor.start(|context| async move {
1054 let cfg = Config {
1056 partition: "test".into(),
1057 codec_config: ((0..).into(), ()),
1058 };
1059 let mut metadata = Metadata::init(context.child("storage"), cfg).await.unwrap();
1060
1061 for i in 0..100 {
1063 metadata.put(U64::new(i), vec![i as u8; 100]);
1064 }
1065
1066 metadata = metadata.sync().await.unwrap();
1070 let buffer = context.encode();
1071 assert!(buffer.contains("sync_rewrites_total 1"), "{buffer}");
1072 assert!(buffer.contains("sync_overwrites_total 0"), "{buffer}");
1073 assert!(
1074 buffer.contains("runtime_storage_write_bytes_total 10912"),
1075 "{buffer}",
1076 );
1077
1078 metadata.put(U64::new(51), vec![0xff; 100]);
1080
1081 metadata = metadata.sync().await.unwrap();
1083 let buffer = context.encode();
1084 assert!(buffer.contains("sync_rewrites_total 2"), "{buffer}");
1085 assert!(buffer.contains("sync_overwrites_total 0"), "{buffer}");
1086 assert!(
1087 buffer.contains("runtime_storage_write_bytes_total 21824"),
1088 "{buffer}",
1089 );
1090
1091 metadata = metadata.sync().await.unwrap();
1095 let buffer = context.encode();
1096 assert!(buffer.contains("sync_rewrites_total 2"), "{buffer}");
1097 assert!(buffer.contains("sync_overwrites_total 1"), "{buffer}");
1098 assert!(
1099 buffer.contains("runtime_storage_write_bytes_total 21937"),
1100 "{buffer}",
1101 );
1102
1103 metadata = metadata.sync().await.unwrap();
1105 let buffer = context.encode();
1106 assert!(buffer.contains("sync_rewrites_total 2"), "{buffer}");
1107 assert!(buffer.contains("sync_overwrites_total 1"), "{buffer}");
1108 assert!(
1109 buffer.contains("runtime_storage_write_bytes_total 21937"),
1110 "{buffer}",
1111 );
1112
1113 metadata.remove(&U64::new(51));
1117 metadata = metadata.sync().await.unwrap();
1118 let buffer = context.encode();
1119 assert!(buffer.contains("sync_rewrites_total 3"), "{buffer}");
1120 assert!(buffer.contains("sync_overwrites_total 1"), "{buffer}");
1121 assert!(
1122 buffer.contains("runtime_storage_write_bytes_total 32740"),
1123 "{buffer}"
1124 );
1125
1126 metadata = metadata.sync().await.unwrap();
1128 let buffer = context.encode();
1129 assert!(buffer.contains("sync_rewrites_total 4"), "{buffer}");
1130 assert!(buffer.contains("sync_overwrites_total 1"), "{buffer}");
1131 assert!(
1132 buffer.contains("runtime_storage_write_bytes_total 43543"),
1133 "{buffer}"
1134 );
1135
1136 metadata.put(U64::new(50), vec![0xff; 100]);
1140 metadata = metadata.sync().await.unwrap();
1141 let buffer = context.encode();
1142 assert!(buffer.contains("sync_rewrites_total 4"), "{buffer}");
1143 assert!(buffer.contains("sync_overwrites_total 2"), "{buffer}");
1144 assert!(
1145 buffer.contains("runtime_storage_write_bytes_total 43656"),
1146 "{buffer}"
1147 );
1148
1149 metadata.destroy().await.unwrap();
1151 });
1152 }
1153
1154 #[test_traced]
1155 fn test_multi_key_overwrites() {
1156 let executor = deterministic::Runner::default();
1157 executor.start(|context| async move {
1158 let cfg = Config {
1159 partition: "test".into(),
1160 codec_config: ((0..).into(), ()),
1161 };
1162 let mut metadata =
1163 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
1164 .await
1165 .unwrap();
1166
1167 for i in 0..100 {
1169 metadata.put(U64::new(i), vec![i as u8; 100]);
1170 }
1171 metadata = metadata.sync().await.unwrap();
1172 metadata = metadata.sync().await.unwrap();
1173 let buffer = context.encode();
1174 assert!(buffer.contains("first_sync_rewrites_total 2"), "{buffer}");
1175 assert!(
1176 buffer.contains("runtime_storage_write_bytes_total 21824"),
1177 "{buffer}",
1178 );
1179
1180 for i in [10u64, 11, 12, 50, 98, 99] {
1182 metadata.put(U64::new(i), vec![0xAA; 100]);
1183 }
1184
1185 metadata = metadata.sync().await.unwrap();
1190 let buffer = context.encode();
1191 assert!(buffer.contains("first_sync_rewrites_total 2"), "{buffer}");
1192 assert!(buffer.contains("first_sync_overwrites_total 1"), "{buffer}");
1193 assert!(
1194 buffer.contains("runtime_storage_write_bytes_total 22442"),
1195 "{buffer}",
1196 );
1197
1198 metadata = metadata.sync().await.unwrap();
1200 let buffer = context.encode();
1201 assert!(buffer.contains("first_sync_overwrites_total 2"), "{buffer}");
1202 assert!(
1203 buffer.contains("runtime_storage_write_bytes_total 23060"),
1204 "{buffer}",
1205 );
1206
1207 metadata = metadata.sync().await.unwrap();
1209 let buffer = context.encode();
1210 assert!(buffer.contains("first_sync_rewrites_total 2"), "{buffer}");
1211 assert!(buffer.contains("first_sync_overwrites_total 2"), "{buffer}");
1212 assert!(
1213 buffer.contains("runtime_storage_write_bytes_total 23060"),
1214 "{buffer}",
1215 );
1216
1217 metadata.put(U64::new(20), vec![0xBB; 100]);
1221 metadata.put(U64::new(30), vec![0xCC; 150]);
1222 metadata = metadata.sync().await.unwrap();
1223 metadata = metadata.sync().await.unwrap();
1224 let buffer = context.encode();
1225 assert!(buffer.contains("first_sync_rewrites_total 4"), "{buffer}");
1226 assert!(buffer.contains("first_sync_overwrites_total 2"), "{buffer}");
1227
1228 drop(metadata);
1230 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
1231 .await
1232 .unwrap();
1233
1234 for i in 0..100u64 {
1236 let expected = match i {
1237 10 | 11 | 12 | 50 | 98 | 99 => vec![0xAA; 100],
1238 20 => vec![0xBB; 100],
1239 30 => vec![0xCC; 150],
1240 _ => vec![i as u8; 100],
1241 };
1242 assert_eq!(metadata.get(&U64::new(i)).unwrap(), &expected, "key {i}");
1243 }
1244
1245 metadata.destroy().await.unwrap();
1246 });
1247 }
1248
1249 #[test_traced]
1250 fn test_sync_with_no_changes() {
1251 let executor = deterministic::Runner::default();
1252 executor.start(|context| async move {
1253 let cfg = Config {
1254 partition: "test".into(),
1255 codec_config: ((0..).into(), ()),
1256 };
1257 let mut metadata =
1258 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
1259 .await
1260 .unwrap();
1261
1262 metadata = metadata
1264 .put_sync(U64::new(1), b"hello".to_vec())
1265 .await
1266 .unwrap();
1267
1268 metadata = metadata.sync().await.unwrap();
1271 let buffer = context.encode();
1272 assert!(buffer.contains("sync_rewrites_total 2"));
1273 assert!(buffer.contains("sync_overwrites_total 0"));
1274
1275 metadata = metadata.sync().await.unwrap();
1277 let buffer = context.encode();
1278 assert!(buffer.contains("sync_rewrites_total 2"));
1279 assert!(buffer.contains("sync_overwrites_total 0"));
1280
1281 metadata = metadata.sync().await.unwrap();
1283 let buffer = context.encode();
1284 assert!(buffer.contains("sync_rewrites_total 2"));
1285 assert!(buffer.contains("sync_overwrites_total 0"));
1286
1287 drop(metadata);
1289 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
1290 .await
1291 .unwrap();
1292 assert_eq!(metadata.get(&U64::new(1)).unwrap(), b"hello");
1293
1294 metadata.destroy().await.unwrap();
1295 });
1296 }
1297
1298 #[test_traced]
1299 fn test_get_mut_marks_modified() {
1300 let executor = deterministic::Runner::default();
1301 executor.start(|context| async move {
1302 let cfg = Config {
1303 partition: "test".into(),
1304 codec_config: ((0..).into(), ()),
1305 };
1306 let mut metadata =
1307 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
1308 .await
1309 .unwrap();
1310
1311 metadata = metadata
1313 .put_sync(U64::new(1), b"hello".to_vec())
1314 .await
1315 .unwrap();
1316
1317 metadata = metadata.sync().await.unwrap();
1319
1320 let value = metadata.get_mut(&U64::new(1)).unwrap();
1322 value[0] = b'H';
1323
1324 metadata = metadata.sync().await.unwrap();
1326 let buffer = context.encode();
1327 assert!(buffer.contains("first_sync_rewrites_total 2"));
1328 assert!(buffer.contains("first_sync_overwrites_total 1"));
1329
1330 drop(metadata);
1332 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
1333 .await
1334 .unwrap();
1335
1336 let value = metadata.get(&U64::new(1)).unwrap();
1338 assert_eq!(value[0], b'H');
1339
1340 metadata.destroy().await.unwrap();
1341 });
1342 }
1343
1344 #[test_traced]
1345 fn test_mixed_operation_sequences() {
1346 let executor = deterministic::Runner::default();
1347 executor.start(|context| async move {
1348 let cfg = Config {
1349 partition: "test".into(),
1350 codec_config: ((0..).into(), ()),
1351 };
1352 let mut metadata =
1353 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
1354 .await
1355 .unwrap();
1356
1357 let key = U64::new(1);
1358
1359 metadata.put(key.clone(), b"first".to_vec());
1361 metadata.remove(&key);
1362 metadata = metadata
1363 .put_sync(key.clone(), b"second".to_vec())
1364 .await
1365 .unwrap();
1366 let value = metadata.get(&key).unwrap();
1367 assert_eq!(value, b"second");
1368
1369 metadata.put(key.clone(), b"third".to_vec());
1371 let value = metadata.get_mut(&key).unwrap();
1372 value[0] = b'T';
1373 metadata.remove(&key);
1374 metadata = metadata
1375 .put_sync(key.clone(), b"fourth".to_vec())
1376 .await
1377 .unwrap();
1378 let value = metadata.get(&key).unwrap();
1379 assert_eq!(value, b"fourth");
1380
1381 drop(metadata);
1383 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
1384 .await
1385 .unwrap();
1386
1387 let value = metadata.get(&key).unwrap();
1389 assert_eq!(value, b"fourth");
1390
1391 metadata.destroy().await.unwrap();
1392 });
1393 }
1394
1395 #[test_traced]
1396 fn test_overwrite_vs_rewrite() {
1397 let executor = deterministic::Runner::default();
1398 executor.start(|context| async move {
1399 let cfg = Config {
1400 partition: "test".into(),
1401 codec_config: ((0..).into(), ()),
1402 };
1403 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("storage"), cfg)
1404 .await
1405 .unwrap();
1406
1407 metadata.put(U64::new(1), vec![1; 10]);
1409 metadata.put(U64::new(2), vec![2; 10]);
1410 metadata = metadata.sync().await.unwrap();
1411
1412 metadata.put(U64::new(1), vec![0xFF; 10]);
1414 metadata = metadata.sync().await.unwrap();
1415 let buffer = context.encode();
1416 assert!(buffer.contains("sync_rewrites_total 2"));
1417 assert!(buffer.contains("sync_overwrites_total 0"));
1418
1419 metadata = metadata.sync().await.unwrap();
1421 let buffer = context.encode();
1422 assert!(buffer.contains("sync_rewrites_total 2"));
1423 assert!(buffer.contains("sync_overwrites_total 1"));
1424
1425 metadata.put(U64::new(1), vec![0xAA; 10]);
1427 metadata = metadata.sync().await.unwrap();
1428 let buffer = context.encode();
1429 assert!(buffer.contains("sync_rewrites_total 2"));
1430 assert!(buffer.contains("sync_overwrites_total 2"));
1431
1432 metadata.put(U64::new(1), vec![0xFF; 20]);
1434 metadata = metadata.sync().await.unwrap();
1435 let buffer = context.encode();
1436 assert!(buffer.contains("sync_rewrites_total 3"));
1437 assert!(buffer.contains("sync_overwrites_total 2"));
1438
1439 metadata.put(U64::new(3), vec![3; 10]);
1441 metadata = metadata.sync().await.unwrap();
1442 let buffer = context.encode();
1443 assert!(buffer.contains("sync_rewrites_total 4"));
1444 assert!(buffer.contains("sync_overwrites_total 2"));
1445
1446 metadata = metadata.sync().await.unwrap();
1448 let buffer = context.encode();
1449 assert!(buffer.contains("sync_rewrites_total 5"));
1450 assert!(buffer.contains("sync_overwrites_total 2"));
1451
1452 metadata.put(U64::new(2), vec![0xAA; 10]);
1454 metadata = metadata.sync().await.unwrap();
1455 let buffer = context.encode();
1456 assert!(buffer.contains("sync_rewrites_total 5"));
1457 assert!(buffer.contains("sync_overwrites_total 3"));
1458
1459 metadata.destroy().await.unwrap();
1460 });
1461 }
1462
1463 #[test_traced]
1464 fn test_blob_resize() {
1465 let executor = deterministic::Runner::default();
1466 executor.start(|context| async move {
1467 let cfg = Config {
1468 partition: "test".into(),
1469 codec_config: ((0..).into(), ()),
1470 };
1471 let mut metadata =
1472 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
1473 .await
1474 .unwrap();
1475
1476 for i in 0..10 {
1478 metadata.put(U64::new(i), vec![i as u8; 100]);
1479 }
1480 metadata = metadata.sync().await.unwrap();
1481
1482 metadata = metadata.sync().await.unwrap();
1484 let buffer = context.encode();
1485 assert!(buffer.contains("first_sync_rewrites_total 2"));
1486 assert!(buffer.contains("first_sync_overwrites_total 0"));
1487
1488 for i in 1..10 {
1490 metadata.remove(&U64::new(i));
1491 }
1492 metadata = metadata.sync().await.unwrap();
1493
1494 let value = metadata.get(&U64::new(0)).unwrap();
1496 assert_eq!(value.len(), 100);
1497 assert_eq!(value[0], 0);
1498
1499 let buffer = context.encode();
1501 assert!(buffer.contains("first_sync_rewrites_total 3"));
1502 assert!(buffer.contains("first_sync_overwrites_total 0"));
1503
1504 drop(metadata);
1506 let metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
1507 .await
1508 .unwrap();
1509
1510 let value = metadata.get(&U64::new(0)).unwrap();
1512 assert_eq!(value.len(), 100);
1513 assert_eq!(value[0], 0);
1514
1515 for i in 1..10 {
1517 assert!(metadata.get(&U64::new(i)).is_none());
1518 }
1519
1520 metadata.destroy().await.unwrap();
1521 });
1522 }
1523
1524 #[test_traced]
1525 fn test_clear_and_repopulate() {
1526 let executor = deterministic::Runner::default();
1527 executor.start(|context| async move {
1528 let cfg = Config {
1529 partition: "test".into(),
1530 codec_config: ((0..).into(), ()),
1531 };
1532 let mut metadata =
1533 Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg.clone())
1534 .await
1535 .unwrap();
1536
1537 metadata.put(U64::new(1), b"first".to_vec());
1539 metadata = metadata
1540 .put_sync(U64::new(2), b"second".to_vec())
1541 .await
1542 .unwrap();
1543
1544 metadata.clear();
1546 metadata = metadata.sync().await.unwrap();
1547
1548 assert!(metadata.get(&U64::new(1)).is_none());
1550 assert!(metadata.get(&U64::new(2)).is_none());
1551
1552 drop(metadata);
1554 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
1555 .await
1556 .unwrap();
1557
1558 assert!(metadata.get(&U64::new(1)).is_none());
1560 assert!(metadata.get(&U64::new(2)).is_none());
1561
1562 metadata.put(U64::new(3), b"third".to_vec());
1564 metadata = metadata
1565 .put_sync(U64::new(4), b"fourth".to_vec())
1566 .await
1567 .unwrap();
1568
1569 assert_eq!(metadata.get(&U64::new(3)).unwrap(), b"third");
1571 assert_eq!(metadata.get(&U64::new(4)).unwrap(), b"fourth");
1572 assert!(metadata.get(&U64::new(1)).is_none());
1573 assert!(metadata.get(&U64::new(2)).is_none());
1574
1575 metadata.destroy().await.unwrap();
1576 });
1577 }
1578
1579 fn test_metadata_operations_and_restart(num_operations: usize) -> String {
1580 let executor = deterministic::Runner::default();
1581 executor.start(|mut context| async move {
1582 let cfg = Config {
1583 partition: "test-determinism".into(),
1584 codec_config: ((0..).into(), ()),
1585 };
1586 let mut metadata =
1587 Metadata::<_, U64, Vec<u8>>::init(context.child("storage"), cfg.clone())
1588 .await
1589 .unwrap();
1590
1591 for i in 0..num_operations {
1593 let key = U64::new(i as u64);
1594 let mut value = vec![0u8; 64];
1595 context.fill_bytes(&mut value);
1596 metadata.put(key, value);
1597
1598 if context.random_bool(0.1) {
1600 metadata = metadata.sync().await.unwrap();
1601 }
1602
1603 if context.random_bool(0.1) {
1605 let selected_index = context.random_range(0..=i);
1606 let update_key = U64::new(selected_index as u64);
1607 let mut new_value = vec![0u8; 64];
1608 context.fill_bytes(&mut new_value);
1609 metadata.put(update_key, new_value);
1610 }
1611
1612 if context.random_bool(0.1) {
1614 let selected_index = context.random_range(0..=i);
1615 let remove_key = U64::new(selected_index as u64);
1616 metadata.remove(&remove_key);
1617 }
1618
1619 if context.random_bool(0.1) {
1621 let selected_index = context.random_range(0..=i);
1622 let mut_key = U64::new(selected_index as u64);
1623 if let Some(value) = metadata.get_mut(&mut_key)
1624 && !value.is_empty()
1625 {
1626 value[0] = value[0].wrapping_add(1);
1627 }
1628 }
1629 }
1630 metadata = metadata.sync().await.unwrap();
1631
1632 metadata.destroy().await.unwrap();
1634
1635 context.auditor().state()
1636 })
1637 }
1638
1639 #[test_group("slow")]
1640 #[test_traced]
1641 fn test_determinism() {
1642 let state1 = test_metadata_operations_and_restart(1_000);
1643 let state2 = test_metadata_operations_and_restart(1_000);
1644 assert_eq!(state1, state2);
1645 }
1646
1647 #[test_traced]
1648 fn test_keys_iterator() {
1649 let executor = deterministic::Runner::default();
1651 executor.start(|context| async move {
1652 let cfg = Config {
1654 partition: "test".into(),
1655 codec_config: ((0..).into(), ()),
1656 };
1657 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("storage"), cfg)
1658 .await
1659 .unwrap();
1660
1661 metadata.put(U64::new(0x1000), b"value1".to_vec());
1663 metadata.put(U64::new(0x1001), b"value2".to_vec());
1664 metadata.put(U64::new(0x1002), b"value3".to_vec());
1665 metadata.put(U64::new(0x2000), b"value4".to_vec());
1666 metadata.put(U64::new(0x2001), b"value5".to_vec());
1667 metadata.put(U64::new(0x3000), b"value6".to_vec());
1668
1669 let all_keys: Vec<_> = metadata.keys().cloned().collect();
1671 assert_eq!(all_keys.len(), 6);
1672 assert!(all_keys.contains(&U64::new(0x1000)));
1673 assert!(all_keys.contains(&U64::new(0x3000)));
1674
1675 let prefix = hex!("0x00000000000010");
1677 let prefix_keys: Vec<_> = metadata
1678 .keys()
1679 .filter(|k| k.as_ref().starts_with(&prefix))
1680 .cloned()
1681 .collect();
1682 assert_eq!(prefix_keys.len(), 3);
1683 assert!(prefix_keys.contains(&U64::new(0x1000)));
1684 assert!(prefix_keys.contains(&U64::new(0x1001)));
1685 assert!(prefix_keys.contains(&U64::new(0x1002)));
1686 assert!(!prefix_keys.contains(&U64::new(0x2000)));
1687
1688 let prefix = hex!("0x00000000000020");
1690 let prefix_keys: Vec<_> = metadata
1691 .keys()
1692 .filter(|k| k.as_ref().starts_with(&prefix))
1693 .cloned()
1694 .collect();
1695 assert_eq!(prefix_keys.len(), 2);
1696 assert!(prefix_keys.contains(&U64::new(0x2000)));
1697 assert!(prefix_keys.contains(&U64::new(0x2001)));
1698
1699 let prefix = hex!("0x00000000000040");
1701 let prefix_keys: Vec<_> = metadata
1702 .keys()
1703 .filter(|k| k.as_ref().starts_with(&prefix))
1704 .cloned()
1705 .collect();
1706 assert_eq!(prefix_keys.len(), 0);
1707
1708 metadata.destroy().await.unwrap();
1709 });
1710 }
1711
1712 #[test_traced]
1713 fn test_retain() {
1714 let executor = deterministic::Runner::default();
1716 executor.start(|context| async move {
1717 let cfg = Config {
1719 partition: "test".into(),
1720 codec_config: ((0..).into(), ()),
1721 };
1722 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("first"), cfg)
1723 .await
1724 .unwrap();
1725
1726 metadata.put(U64::new(0x1000), b"value1".to_vec());
1728 metadata.put(U64::new(0x1001), b"value2".to_vec());
1729 metadata.put(U64::new(0x1002), b"value3".to_vec());
1730 metadata.put(U64::new(0x2000), b"value4".to_vec());
1731 metadata.put(U64::new(0x2001), b"value5".to_vec());
1732 metadata.put(U64::new(0x3000), b"value6".to_vec());
1733
1734 let buffer = context.encode();
1736 assert!(buffer.contains("first_keys 6"));
1737
1738 let prefix = hex!("0x00000000000010");
1740 metadata.retain(|k, _| !k.as_ref().starts_with(&prefix));
1741
1742 let buffer = context.encode();
1744 assert!(buffer.contains("first_keys 3"));
1745
1746 assert!(metadata.get(&U64::new(0x1000)).is_none());
1748 assert!(metadata.get(&U64::new(0x1001)).is_none());
1749 assert!(metadata.get(&U64::new(0x1002)).is_none());
1750 assert!(metadata.get(&U64::new(0x2000)).is_some());
1751 assert!(metadata.get(&U64::new(0x2001)).is_some());
1752 assert!(metadata.get(&U64::new(0x3000)).is_some());
1753
1754 metadata.sync().await.unwrap();
1756 let cfg = Config {
1757 partition: "test".into(),
1758 codec_config: ((0..).into(), ()),
1759 };
1760 let mut metadata = Metadata::<_, U64, Vec<u8>>::init(context.child("second"), cfg)
1761 .await
1762 .unwrap();
1763
1764 assert!(metadata.get(&U64::new(0x1000)).is_none());
1766 assert!(metadata.get(&U64::new(0x2000)).is_some());
1767 assert_eq!(metadata.keys().count(), 3);
1768
1769 let prefix = hex!("0x00000000000040");
1771 metadata.retain(|k, _| !k.as_ref().starts_with(&prefix));
1772
1773 metadata.retain(|_, _| false);
1775 assert_eq!(metadata.keys().count(), 0);
1776
1777 metadata.destroy().await.unwrap();
1778 });
1779 }
1780}