1use core::future::{Future, poll_fn};
19use core::pin::{Pin, pin};
20use core::task::{Context, Poll, ready};
21use std::collections::VecDeque;
22
23use bytes::{BufMut, Bytes, BytesMut};
24use futures_core::Stream;
25use mkit_core::hash::Hash;
26
27use super::codec;
28use super::error::StoreError;
29use super::keys;
30use super::kv::{
31 Batch, BatchOutcome, Cursor, Key, KeyClasses, MAX_BATCH_BYTES, MAX_BATCH_OPS, MAX_KEY_BYTES,
32 MAX_VALUE_BYTES, NamespaceStore, Precondition, Value,
33};
34use super::partition::Partition;
35use crate::rt::{BoxFuture, MaybeSend, MaybeSync};
36
37pub const EXPORT_MAGIC: [u8; 8] = *b"mkitexp\0";
39pub const EXPORT_FORMAT_V1: u8 = 1;
41pub const EXPORT_END: [u8; 2] = [0, 0];
44
45const HEADER_LEN: usize = EXPORT_MAGIC.len() + 1 + 4 + 8;
46const EXPORT_PAGE: u32 = 256;
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51#[non_exhaustive]
52pub struct ExportHeader {
53 pub layout_version: u32,
58 pub exported_at_ms: u64,
60}
61
62impl ExportHeader {
63 #[must_use]
65 pub fn new(layout_version: u32, exported_at_ms: u64) -> Self {
66 Self {
67 layout_version,
68 exported_at_ms,
69 }
70 }
71}
72
73#[derive(Debug, Clone, PartialEq, Eq)]
75#[non_exhaustive]
76pub struct ExportRecord {
77 pub partition: Partition,
79 pub key: Key,
81 pub value: Value,
83}
84
85impl ExportRecord {
86 #[must_use]
88 pub fn new(partition: Partition, key: Key, value: Value) -> Self {
89 Self {
90 partition,
91 key,
92 value,
93 }
94 }
95}
96
97#[derive(Debug, Clone, PartialEq, Eq, Default)]
99#[non_exhaustive]
100pub struct ExportPage {
101 pub records: Vec<ExportRecord>,
103 pub next: Option<Cursor>,
105}
106
107fn export_range<S: NamespaceStore>(store: &S) -> (Key, Key) {
110 if store.capabilities().key_classes == KeyClasses::RefsOnly {
111 keys::class_range(keys::TAG_REF)
112 } else {
113 (Key::default(), Key::new(vec![0xff]))
114 }
115}
116
117pub async fn export_header<S: NamespaceStore>(
119 store: &S,
120 p: &Partition,
121 now_ms: u64,
122) -> Result<ExportHeader, StoreError> {
123 let layout_version = match store.capabilities().implicit_layout_version {
124 Some(version) => version,
125 None => match store.get(p, &keys::layout_version()).await? {
126 Some(value) => codec::decode_u32(&value)?,
127 None => keys::LAYOUT_VERSION,
128 },
129 };
130 Ok(ExportHeader {
131 layout_version,
132 exported_at_ms: now_ms,
133 })
134}
135
136pub async fn export_page<S: NamespaceStore>(
139 store: &S,
140 p: &Partition,
141 after: Option<&Cursor>,
142 limit: u32,
143) -> Result<ExportPage, StoreError> {
144 let (start, end) = export_range(store);
145 let page = store.scan(p, &start, &end, after, limit).await?;
146 Ok(ExportPage {
147 records: page
148 .entries
149 .into_iter()
150 .map(|(key, value)| ExportRecord {
151 partition: p.clone(),
152 key,
153 value,
154 })
155 .collect(),
156 next: page.next,
157 })
158}
159
160async fn owned_page<S: NamespaceStore>(
161 store: &S,
162 p: Partition,
163 after: Option<Cursor>,
164) -> Result<ExportPage, StoreError> {
165 export_page(store, &p, after.as_ref(), EXPORT_PAGE).await
166}
167
168pub async fn export_partition<'a, S: NamespaceStore>(
172 store: &'a S,
173 p: &Partition,
174 now_ms: u64,
175) -> Result<(ExportHeader, ExportStream<'a, S>), StoreError> {
176 let header = export_header(store, p, now_ms).await?;
177 let stream = ExportStream {
178 store,
179 partition: p.clone(),
180 pending: Some(Box::pin(owned_page(store, p.clone(), None))),
181 buffered: VecDeque::new(),
182 };
183 Ok((header, stream))
184}
185
186pub struct ExportStream<'a, S> {
188 store: &'a S,
189 partition: Partition,
190 pending: Option<BoxFuture<'a, Result<ExportPage, StoreError>>>,
191 buffered: VecDeque<ExportRecord>,
192}
193
194impl<S> core::fmt::Debug for ExportStream<'_, S> {
195 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
196 f.debug_struct("ExportStream")
197 .field("partition", &self.partition)
198 .field("buffered", &self.buffered.len())
199 .finish_non_exhaustive()
200 }
201}
202
203impl<S: NamespaceStore> Stream for ExportStream<'_, S> {
204 type Item = Result<ExportRecord, StoreError>;
205
206 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
207 let this = self.get_mut();
208 loop {
209 if let Some(record) = this.buffered.pop_front() {
210 return Poll::Ready(Some(Ok(record)));
211 }
212 let Some(pending) = this.pending.as_mut() else {
213 return Poll::Ready(None);
214 };
215 let page = ready!(pending.as_mut().poll(cx));
216 this.pending = None;
217 let page = match page {
218 Ok(page) => page,
219 Err(e) => return Poll::Ready(Some(Err(e))),
220 };
221 this.buffered.extend(page.records);
222 if let Some(next) = page.next {
223 let fut = owned_page(this.store, this.partition.clone(), Some(next));
224 this.pending = Some(Box::pin(fut));
225 }
226 }
227 }
228}
229
230#[must_use]
233pub fn encode_export_header(header: &ExportHeader) -> Bytes {
234 let mut buf = BytesMut::with_capacity(HEADER_LEN);
235 buf.put_slice(&EXPORT_MAGIC);
236 buf.put_u8(EXPORT_FORMAT_V1);
237 buf.put_u32(header.layout_version);
238 buf.put_u64(header.exported_at_ms);
239 buf.freeze()
240}
241
242pub fn encode_export_record(record: &ExportRecord) -> Result<Bytes, StoreError> {
251 let partition = record.partition.encode()?;
252 let (key, value) = (record.key.as_bytes(), record.value.as_bytes());
253 let too_long = || StoreError::Invalid("export record exceeds limits".into());
254 if key.len() > MAX_KEY_BYTES || value.len() > MAX_VALUE_BYTES {
255 return Err(too_long());
256 }
257 let mut buf = BytesMut::with_capacity(8 + partition.len() + key.len() + value.len());
258 buf.put_u16(u16::try_from(partition.len()).map_err(|_| too_long())?);
259 buf.put_slice(&partition);
260 buf.put_u16(u16::try_from(key.len()).map_err(|_| too_long())?);
261 buf.put_slice(key);
262 buf.put_u32(u32::try_from(value.len()).map_err(|_| too_long())?);
263 buf.put_slice(value);
264 Ok(buf.freeze())
265}
266
267#[derive(Debug)]
269pub struct ExportReader<'a> {
270 rest: &'a [u8],
271 done: bool,
272}
273
274fn corrupt(what: &'static str) -> StoreError {
275 StoreError::Corrupt(what.into())
276}
277
278fn take<'a>(rest: &mut &'a [u8], n: usize) -> Result<&'a [u8], StoreError> {
279 if rest.len() < n {
280 return Err(corrupt("truncated export"));
281 }
282 let (head, tail) = rest.split_at(n);
283 *rest = tail;
284 Ok(head)
285}
286
287fn take_len(rest: &mut &[u8], width: usize, max: usize) -> Result<usize, StoreError> {
288 let len = take(rest, width)?
289 .iter()
290 .fold(0_usize, |n, &b| (n << 8) | usize::from(b));
291 if len > max {
292 return Err(corrupt("export record exceeds limits"));
293 }
294 Ok(len)
295}
296
297impl<'a> ExportReader<'a> {
298 pub fn new(bytes: &'a [u8]) -> Result<(ExportHeader, Self), StoreError> {
305 let (head, rest) = bytes
306 .split_first_chunk::<HEADER_LEN>()
307 .ok_or_else(|| corrupt("truncated export"))?;
308 let (magic, fields) = head.split_at(EXPORT_MAGIC.len());
309 if magic != EXPORT_MAGIC {
310 return Err(corrupt("not an mkit export"));
311 }
312 if fields[0] != EXPORT_FORMAT_V1 {
313 return Err(corrupt("unknown export format"));
314 }
315 let mut version = [0; 4];
316 version.copy_from_slice(&fields[1..5]);
317 let mut at = [0; 8];
318 at.copy_from_slice(&fields[5..]);
319 let header = ExportHeader {
320 layout_version: u32::from_be_bytes(version),
321 exported_at_ms: u64::from_be_bytes(at),
322 };
323 Ok((header, Self { rest, done: false }))
324 }
325
326 fn record(&mut self) -> Result<Option<ExportRecord>, StoreError> {
327 let partition_len = take_len(&mut self.rest, 2, usize::from(u16::MAX))?;
328 if partition_len == 0 {
329 if !self.rest.is_empty() {
330 return Err(corrupt("bytes after the export end marker"));
331 }
332 return Ok(None);
333 }
334 let partition = Partition::decode(take(&mut self.rest, partition_len)?)?;
335 let key_len = take_len(&mut self.rest, 2, MAX_KEY_BYTES)?;
336 let key = Key::new(take(&mut self.rest, key_len)?.to_vec());
337 let value_len = take_len(&mut self.rest, 4, MAX_VALUE_BYTES)?;
338 let value = Value::new(take(&mut self.rest, value_len)?.to_vec());
339 Ok(Some(ExportRecord {
340 partition,
341 key,
342 value,
343 }))
344 }
345}
346
347impl Iterator for ExportReader<'_> {
348 type Item = Result<ExportRecord, StoreError>;
349
350 fn next(&mut self) -> Option<Self::Item> {
351 if self.done {
352 return None;
353 }
354 let record = self.record();
355 self.done = !matches!(record, Ok(Some(_)));
356 record.transpose()
357 }
358}
359
360#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
362pub enum ImportMode {
363 #[default]
365 Fresh,
366 Merge,
368}
369
370#[derive(Debug)]
383pub struct Importer<'a, S> {
384 store: &'a S,
385 mode: ImportMode,
386 layout_version: u32,
387 holds_v: bool,
388 max_ops: usize,
389 partition: Option<Partition>,
390 saw_v: bool,
391 batch: Batch,
392 bytes: usize,
393 imported: u64,
394}
395
396fn check_layout(version: u32) -> Result<(), StoreError> {
397 if version > keys::LAYOUT_VERSION {
398 return Err(StoreError::Unsupported(
399 "export has a newer layout version than this binary".into(),
400 ));
401 }
402 Ok(())
403}
404
405impl<'a, S: NamespaceStore> Importer<'a, S> {
406 pub fn new(store: &'a S, header: &ExportHeader, mode: ImportMode) -> Result<Self, StoreError> {
413 check_layout(header.layout_version)?;
414 let caps = store.capabilities();
415 if caps
416 .implicit_layout_version
417 .is_some_and(|implicit| implicit != header.layout_version)
418 {
419 return Err(StoreError::Unsupported(
420 "export layout version differs from the store's implicit version".into(),
421 ));
422 }
423 Ok(Self {
424 store,
425 mode,
426 layout_version: header.layout_version,
427 holds_v: caps.key_classes == KeyClasses::All,
428 max_ops: if caps.atomic_multi_key {
429 MAX_BATCH_OPS
430 } else {
431 1
432 },
433 partition: None,
434 saw_v: false,
435 batch: Batch::new(),
436 bytes: 0,
437 imported: 0,
438 })
439 }
440
441 pub async fn push(&mut self, record: ExportRecord) -> Result<(), StoreError> {
449 let is_v = record.key == keys::layout_version();
450 if is_v {
451 check_layout(codec::decode_u32(&record.value)?)?;
452 }
453 let size = record.key.as_bytes().len() + record.value.as_bytes().len();
454 if self.partition.as_ref() != Some(&record.partition) {
455 self.end_partition().await?;
456 if self.mode == ImportMode::Fresh {
457 let (start, end) = export_range(self.store);
458 let page = self
459 .store
460 .scan(&record.partition, &start, &end, None, 1)
461 .await?;
462 if !page.entries.is_empty() {
463 return Err(StoreError::Invalid(
464 "import target partition is not empty".into(),
465 ));
466 }
467 }
468 self.partition = Some(record.partition);
469 self.saw_v = false;
470 } else if self.batch.writes.len() >= self.max_ops || self.bytes + size > MAX_BATCH_BYTES {
471 self.flush().await?;
472 }
473 self.saw_v |= is_v;
474 self.batch = core::mem::take(&mut self.batch).put(record.key, record.value);
475 self.bytes += size;
476 self.imported += 1;
477 Ok(())
478 }
479
480 async fn apply(&self, p: &Partition, batch: Batch) -> Result<(), StoreError> {
481 match self.store.apply(p, batch).await? {
482 BatchOutcome::Committed => Ok(()),
483 _ => Err(StoreError::unavailable("import batch did not commit")),
484 }
485 }
486
487 async fn flush(&mut self) -> Result<(), StoreError> {
488 let batch = core::mem::take(&mut self.batch);
489 self.bytes = 0;
490 match &self.partition {
491 Some(p) if !batch.writes.is_empty() => self.apply(p, batch).await,
492 _ => Ok(()),
493 }
494 }
495
496 async fn end_partition(&mut self) -> Result<(), StoreError> {
499 self.flush().await?;
500 let Some(p) = self.partition.take() else {
501 return Ok(());
502 };
503 if !self.holds_v || self.saw_v {
504 return Ok(());
505 }
506 let key = keys::layout_version();
507 let batch = Batch::new()
508 .require(Precondition::Absent(key.clone()))
509 .put(key, codec::encode_u32(self.layout_version));
510 match self.store.apply(&p, batch).await? {
511 BatchOutcome::Committed | BatchOutcome::PreconditionFailed { .. } => Ok(()),
512 BatchOutcome::DeadlinePassed { .. } => {
513 Err(StoreError::unavailable("import batch did not commit"))
514 }
515 }
516 }
517
518 pub async fn finish(mut self) -> Result<u64, StoreError> {
520 self.end_partition().await?;
521 Ok(self.imported)
522 }
523}
524
525pub async fn import_stream<S, R>(
529 store: &S,
530 header: &ExportHeader,
531 mode: ImportMode,
532 records: R,
533) -> Result<u64, StoreError>
534where
535 S: NamespaceStore,
536 R: Stream<Item = Result<ExportRecord, StoreError>>,
537{
538 let mut importer = Importer::new(store, header, mode)?;
539 let mut records = pin!(records);
540 while let Some(record) = poll_fn(|cx| records.as_mut().poll_next(cx)).await {
541 importer.push(record?).await?;
542 }
543 importer.finish().await
544}
545
546pub trait StoreMaintenance: MaybeSend + MaybeSync {
551 fn layout_version(&self) -> u32;
554
555 fn migrate(&self) -> impl Future<Output = Result<u32, StoreError>> + MaybeSend;
558
559 fn backup_to(&self, dest: &str) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
562}
563
564pub trait StateCommitment: NamespaceStore {
567 fn root(
570 &self,
571 p: &Partition,
572 ) -> impl Future<Output = Result<Option<(String, Hash)>, StoreError>> + MaybeSend;
573
574 fn prove(
577 &self,
578 p: &Partition,
579 key: &Key,
580 ) -> impl Future<Output = Result<Option<Bytes>, StoreError>> + MaybeSend;
581}
582
583#[cfg(test)]
584mod tests {
585 use std::sync::Mutex;
586
587 use futures_executor::block_on;
588
589 use super::*;
590 use crate::memory::MemoryKv;
591 use crate::repo::{NamespaceKey, RepoName};
592 use crate::store::content_index::{
593 BlockEntry, ContentIndex, HoldOutcome, Holder, content_shard,
594 };
595 use crate::store::{PartitionStats, ScanPage, StoreCapabilities};
596
597 fn ns() -> Partition {
598 Partition::Namespace(NamespaceKey::deployment_default())
599 }
600
601 fn repo() -> RepoName {
602 RepoName::new("r").unwrap()
603 }
604
605 fn scan_all<S: NamespaceStore>(store: &S, p: &Partition) -> Vec<ExportRecord> {
607 let mut out = Vec::new();
608 let mut after = None;
609 loop {
610 let page = block_on(export_page(store, p, after.as_ref(), 2)).unwrap();
611 out.extend(page.records);
612 match page.next {
613 Some(next) => after = Some(next),
614 None => return out,
615 }
616 }
617 }
618
619 fn export_bytes<S: NamespaceStore>(store: &S, p: &Partition) -> Vec<u8> {
621 block_on(async {
622 let (header, stream) = export_partition(store, p, 42).await.unwrap();
623 let mut stream = pin!(stream);
624 let mut out = encode_export_header(&header).to_vec();
625 while let Some(record) = poll_fn(|cx| stream.as_mut().poll_next(cx)).await {
626 out.extend_from_slice(&encode_export_record(&record.unwrap()).unwrap());
627 }
628 out.extend_from_slice(&EXPORT_END);
629 out
630 })
631 }
632
633 fn import_bytes<S: NamespaceStore>(
634 store: &S,
635 bytes: &[u8],
636 mode: ImportMode,
637 ) -> Result<u64, StoreError> {
638 let (header, reader) = ExportReader::new(bytes)?;
639 block_on(async {
640 let mut importer = Importer::new(store, &header, mode)?;
641 for record in reader {
642 importer.push(record?).await?;
643 }
644 importer.finish().await
645 })
646 }
647
648 #[test]
649 fn export_import_roundtrip_is_identical() {
650 let (a, b) = ([0x10; 32], [0xf0; 32]);
651 let clock = std::sync::Arc::new(crate::ManualClock::new(0));
652 let idx = ContentIndex::new(MemoryKv::with_clock(clock));
653 let holder = Holder::new(NamespaceKey::deployment_default(), repo());
654 block_on(async {
655 idx.add_holder(&a, &holder, &[9; 32], None, 1)
656 .await
657 .unwrap();
658 let held = idx.add_hold(&b, &[1; 32], 99, 2).await.unwrap();
659 assert_eq!(held, HoldOutcome::Held);
660 let entry = BlockEntry::new("r", 3);
661 idx.block(&b, &entry, 3).await.unwrap();
662 let mut batch = Batch::new()
663 .put(keys::layout_version(), codec::encode_u32(1))
664 .put(keys::grant_epoch(), codec::encode_u64(4));
665 for i in 0..7_u8 {
666 let name = format!("refs/heads/b{i}");
667 let id = codec::encode_ref_id(&[i; 32]);
668 batch = batch.put(keys::ref_key(&repo(), &name), id);
669 }
670 idx.store().apply(&ns(), batch).await.unwrap();
671 });
672 let src = idx.store();
673 let dst = MemoryKv::default();
674 let parts = [ns(), content_shard(&a), content_shard(&b)];
675 for p in &parts {
676 let bytes = export_bytes(src, p);
677 let rows = u64::try_from(scan_all(src, p).len()).unwrap();
678 assert!(rows >= 3);
679 assert_eq!(import_bytes(&dst, &bytes, ImportMode::Fresh).unwrap(), rows);
680 assert_eq!(scan_all(&dst, p), scan_all(src, p), "{p:?}");
681 assert_eq!(export_bytes(&dst, p), bytes, "byte-for-byte {p:?}");
682 }
683 let dst = MemoryKv::default();
685 block_on(async {
686 for p in &parts {
687 let (header, stream) = export_partition(src, p, 42).await.unwrap();
688 import_stream(&dst, &header, ImportMode::Fresh, stream)
689 .await
690 .unwrap();
691 }
692 });
693 for p in &parts {
694 assert_eq!(scan_all(&dst, p), scan_all(src, p));
695 }
696 }
697
698 #[test]
699 fn export_format_golden_bytes() {
700 let header = ExportHeader {
701 layout_version: 1,
702 exported_at_ms: 0x0102_0304_0506_0708,
703 };
704 let record = ExportRecord {
705 partition: Partition::ContentShard(7),
706 key: Key::new(&b"b\0"[..]),
707 value: Value::new(&b"v"[..]),
708 };
709 let bytes = [
710 &encode_export_header(&header)[..],
711 &encode_export_record(&record).unwrap(),
712 &EXPORT_END,
713 ]
714 .concat();
715 let golden: &[u8] = b"mkitexp\0\x01\0\0\0\x01\x01\x02\x03\x04\x05\x06\x07\x08\
716 \0\x03s7\0\0\x02b\0\0\0\0\x01v\0\0";
717 assert_eq!(bytes, golden);
718 let (read, reader) = ExportReader::new(golden).unwrap();
719 assert_eq!(read, header);
720 assert_eq!(
721 reader.collect::<Result<Vec<_>, _>>().unwrap(),
722 vec![record.clone()]
723 );
724 let mut trailing = golden.to_vec();
725 trailing.push(0);
726 let mut bad_magic = golden.to_vec();
727 bad_magic[0] = b'M';
728 let mut newer_format = golden.to_vec();
729 newer_format[8] = 2;
730 let mut long_key = golden[..HEADER_LEN + 5].to_vec();
731 long_key.extend_from_slice(&[0x04, 0x01]);
732 for bad in [
733 &golden[..golden.len() - 1],
734 &golden[..HEADER_LEN],
735 &golden[..5],
736 trailing.as_slice(),
737 bad_magic.as_slice(),
738 newer_format.as_slice(),
739 long_key.as_slice(),
740 ] {
741 let result = ExportReader::new(bad)
742 .and_then(|(_, reader)| reader.collect::<Result<Vec<_>, _>>());
743 assert!(matches!(result, Err(StoreError::Corrupt(_))), "{bad:?}");
744 }
745 let huge = ExportRecord {
746 value: Value::new(vec![0; MAX_VALUE_BYTES + 1]),
747 ..record
748 };
749 assert!(matches!(
750 encode_export_record(&huge),
751 Err(StoreError::Invalid(_))
752 ));
753 }
754
755 #[test]
756 fn import_refuses_newer_layout_version() {
757 let kv = MemoryKv::default();
758 let newer = ExportHeader {
759 layout_version: keys::LAYOUT_VERSION + 1,
760 exported_at_ms: 0,
761 };
762 assert!(matches!(
763 Importer::new(&kv, &newer, ImportMode::Fresh),
764 Err(StoreError::Unsupported(_))
765 ));
766 let bytes = [&encode_export_header(&newer)[..], &EXPORT_END].concat();
767 assert!(matches!(
768 import_bytes(&kv, &bytes, ImportMode::Fresh),
769 Err(StoreError::Unsupported(_))
770 ));
771 let current = ExportHeader {
773 layout_version: keys::LAYOUT_VERSION,
774 exported_at_ms: 0,
775 };
776 let row = ExportRecord {
777 partition: ns(),
778 key: keys::layout_version(),
779 value: codec::encode_u32(keys::LAYOUT_VERSION + 1),
780 };
781 let mut importer = Importer::new(&kv, ¤t, ImportMode::Fresh).unwrap();
782 assert!(matches!(
783 block_on(importer.push(row)),
784 Err(StoreError::Unsupported(_))
785 ));
786 assert!(scan_all(&kv, &ns()).is_empty(), "nothing was written");
787 }
788
789 struct Recording {
791 inner: MemoryKv,
792 batches: Mutex<Vec<(usize, usize)>>,
793 }
794
795 impl NamespaceStore for Recording {
796 fn capabilities(&self) -> StoreCapabilities {
797 self.inner.capabilities()
798 }
799 async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
800 self.inner.get(p, key).await
801 }
802 async fn scan(
803 &self,
804 p: &Partition,
805 start: &Key,
806 end: &Key,
807 after: Option<&Cursor>,
808 limit: u32,
809 ) -> Result<ScanPage, StoreError> {
810 self.inner.scan(p, start, end, after, limit).await
811 }
812 async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
813 let bytes = batch
814 .writes
815 .iter()
816 .map(|w| match w {
817 crate::store::Write::Put(k, v) => k.as_bytes().len() + v.as_bytes().len(),
818 crate::store::Write::Delete(k) => k.as_bytes().len(),
819 })
820 .sum();
821 self.batches
822 .lock()
823 .unwrap()
824 .push((batch.writes.len(), bytes));
825 self.inner.apply(p, batch).await
826 }
827 async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
828 self.inner.stats(p).await
829 }
830 async fn probe(&self) -> Result<(), StoreError> {
831 Ok(())
832 }
833 }
834
835 #[test]
836 fn import_batches_stay_within_do_limits() {
837 let header = ExportHeader {
838 layout_version: keys::LAYOUT_VERSION,
839 exported_at_ms: 0,
840 };
841 let records: Vec<_> = (0..250_u32)
842 .map(|i| {
843 let big = i % 50 < 7;
844 ExportRecord {
845 partition: if i < 200 {
846 ns()
847 } else {
848 content_shard(&[0; 32])
849 },
850 key: keys::ref_key(&repo(), &format!("refs/heads/{i:04}")),
851 value: Value::new(vec![1; if big { 300 * 1024 } else { 32 }]),
852 }
853 })
854 .collect();
855 for caps in [StoreCapabilities::full(), StoreCapabilities::refs_only()] {
856 let store = Recording {
857 inner: MemoryKv::new(caps),
858 batches: Mutex::default(),
859 };
860 let imported = block_on(async {
861 let mut importer = Importer::new(&store, &header, ImportMode::Fresh).unwrap();
862 for record in records.clone() {
863 importer.push(record).await.unwrap();
864 }
865 importer.finish().await.unwrap()
866 });
867 assert_eq!(imported, 250);
868 let batches = store.batches.into_inner().unwrap();
869 let max_ops = if caps.atomic_multi_key {
870 MAX_BATCH_OPS
871 } else {
872 1
873 };
874 for &(ops, bytes) in &batches {
875 assert!(
876 ops <= max_ops && bytes <= MAX_BATCH_BYTES,
877 "{ops} ops, {bytes} B"
878 );
879 }
880 let v_rows = if caps.atomic_multi_key { 2 } else { 0 };
883 assert_eq!(batches.iter().map(|b| b.0).sum::<usize>(), 250 + v_rows);
884 if caps.atomic_multi_key {
885 assert!(batches.len() > 3 && batches.len() < 30, "{batches:?}");
888 }
889 for p in [ns(), content_shard(&[0; 32])] {
890 let want: Vec<_> = records.iter().filter(|r| r.partition == p).collect();
891 let got = scan_all(&store.inner, &p);
892 let got: Vec<_> = got
893 .iter()
894 .filter(|r| r.key != keys::layout_version())
895 .collect();
896 assert_eq!(got, want);
897 }
898 }
899 }
900
901 #[test]
902 fn refs_only_export_is_the_ref_class_with_the_implicit_version() {
903 let kv = MemoryKv::new(StoreCapabilities::refs_only());
904 let key = keys::ref_key(&repo(), "refs/heads/main");
905 block_on(kv.apply(
906 &ns(),
907 Batch::new().put(key.clone(), codec::encode_ref_id(&[5; 32])),
908 ))
909 .unwrap();
910 let header = block_on(export_header(&kv, &ns(), 7)).unwrap();
911 assert_eq!(header.layout_version, keys::LAYOUT_VERSION);
912 let records = scan_all(&kv, &ns());
913 assert_eq!(records.len(), 1);
914 assert_eq!(records[0].key, key);
915 assert_eq!(keys::class_range(keys::TAG_REF), export_range(&kv));
916 let full = MemoryKv::default();
918 let bytes = export_bytes(&kv, &ns());
919 assert_eq!(import_bytes(&full, &bytes, ImportMode::Fresh).unwrap(), 1);
920 let v = block_on(full.get(&ns(), &keys::layout_version())).unwrap();
921 assert_eq!(v, Some(codec::encode_u32(header.layout_version)));
922 assert_eq!(scan_all(&full, &ns()).len(), 2);
923 }
924
925 #[test]
926 fn import_refuses_a_non_empty_target_unless_merging() {
927 let src = MemoryKv::default();
928 let rows = Batch::new()
929 .put(keys::layout_version(), codec::encode_u32(1))
930 .put(keys::grant_epoch(), codec::encode_u64(9));
931 block_on(src.apply(&ns(), rows)).unwrap();
932 let bytes = export_bytes(&src, &ns());
933 let dst = MemoryKv::default();
934 let other = keys::ref_key(&repo(), "refs/heads/keep");
935 let existing = Batch::new()
936 .put(keys::grant_epoch(), codec::encode_u64(1))
937 .put(other.clone(), codec::encode_ref_id(&[1; 32]));
938 block_on(dst.apply(&ns(), existing)).unwrap();
939 let before = scan_all(&dst, &ns());
940 assert!(matches!(
941 import_bytes(&dst, &bytes, ImportMode::Fresh),
942 Err(StoreError::Invalid(_))
943 ));
944 assert_eq!(scan_all(&dst, &ns()), before, "nothing was written");
945 assert_eq!(import_bytes(&dst, &bytes, ImportMode::Merge).unwrap(), 2);
946 let epoch = block_on(dst.get(&ns(), &keys::grant_epoch())).unwrap();
947 assert_eq!(epoch, Some(codec::encode_u64(9)), "imported rows win");
948 let kept = block_on(dst.get(&ns(), &other)).unwrap();
949 assert!(kept.is_some(), "other rows stay");
950 }
951}