Skip to main content

mkit_server/store/
maintenance.rs

1//! Backup and restore (PRD ยง5.3: "a backup and restore procedure and
2//! versioned schema migrations are required for every backend") and the
3//! optional backend hooks.
4//!
5//! **Portable logical export** (R-20) works on every [`NamespaceStore`]: a
6//! full ordered scan of one partition as [`ExportRecord`]s under an
7//! [`ExportHeader`], written in the byte format of [`encode_export_header`]
8//! and [`encode_export_record`] and read back by [`ExportReader`]. A store
9//! never lists its partitions: a full backup exports each partition the
10//! `Partition` enumeration names (namespaces from configuration and the
11//! `nl` list, repos from each coordinator's `rr` registry, index and
12//! content shards by construction, see [`super::content_shards`], ref
13//! shards from the ref-name index). There is no registry of all shards.
14//!
15//! [`StoreMaintenance`] (R-19) and [`StateCommitment`] (R-21) are optional:
16//! nothing in the pipeline requires them.
17
18use 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
37/// First bytes of every export.
38pub const EXPORT_MAGIC: [u8; 8] = *b"mkitexp\0";
39/// The export byte format this binary writes and reads.
40pub const EXPORT_FORMAT_V1: u8 = 1;
41/// End marker: a record whose partition length is 0. A reader requires it,
42/// so a truncated export never imports as a shorter one.
43pub const EXPORT_END: [u8; 2] = [0, 0];
44
45const HEADER_LEN: usize = EXPORT_MAGIC.len() + 1 + 4 + 8;
46/// Records per scan while exporting.
47const EXPORT_PAGE: u32 = 256;
48
49/// What an export was taken from.
50#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51#[non_exhaustive]
52pub struct ExportHeader {
53    /// The partition's key-layout version: its `v` row, a `RefsOnly`
54    /// store's `implicit_layout_version`, or this binary's version for an
55    /// unversioned partition. A dump of several partitions carries the
56    /// highest.
57    pub layout_version: u32,
58    /// When the export started, Unix ms.
59    pub exported_at_ms: u64,
60}
61
62impl ExportHeader {
63    /// A header.
64    #[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/// One exported row.
74#[derive(Debug, Clone, PartialEq, Eq)]
75#[non_exhaustive]
76pub struct ExportRecord {
77    /// The row's partition.
78    pub partition: Partition,
79    /// Key.
80    pub key: Key,
81    /// Value.
82    pub value: Value,
83}
84
85impl ExportRecord {
86    /// A record.
87    #[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/// One page of [`export_page`].
98#[derive(Debug, Clone, PartialEq, Eq, Default)]
99#[non_exhaustive]
100pub struct ExportPage {
101    /// Records in key order.
102    pub records: Vec<ExportRecord>,
103    /// Resume point, if more records may follow.
104    pub next: Option<Cursor>,
105}
106
107/// The key range an export scans: every key of the store's classes. Every
108/// key starts with an ASCII class tag, so `[0xff]` bounds them all.
109fn 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
117/// The header of an export of `p` taken at `now_ms`.
118pub 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
136/// Up to `limit` rows of `p` after `after`, in key order: one stateless
137/// step of an export (a Durable Object serves it per call).
138pub 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
168/// Portable logical backup of one partition: its header and a stream of
169/// every row in key order. The stream is not a snapshot: rows written while
170/// it runs may or may not appear, so export a quiesced partition.
171pub 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
186/// The record stream of [`export_partition`].
187pub 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/// The export header bytes: [`EXPORT_MAGIC`], [`EXPORT_FORMAT_V1`], the
231/// layout version (be32) and the export time (be64).
232#[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
242/// One record's bytes: partition length (be16, never 0) and
243/// [`Partition::encode`] bytes, key length (be16) and key, value length
244/// (be32) and value. An export is the header, its records, then
245/// [`EXPORT_END`].
246///
247/// # Errors
248/// [`StoreError::Invalid`] for an unencodable partition or a key or value
249/// over the contract limits.
250pub 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/// Reads an export's records from its bytes; see [`encode_export_record`].
268#[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    /// Check the header of `bytes` and return it with a reader over the
299    /// records.
300    ///
301    /// # Errors
302    /// [`StoreError::Corrupt`] for a missing magic, an unknown format or a
303    /// short header.
304    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/// How an [`Importer`] treats a partition that already holds rows.
361#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
362pub enum ImportMode {
363    /// Refuse it: a restore goes into empty partitions.
364    #[default]
365    Fresh,
366    /// Write over it: imported rows replace same-key rows, other rows stay.
367    Merge,
368}
369
370/// Restores exported records into a store, in batches within
371/// [`MAX_BATCH_OPS`] and [`MAX_BATCH_BYTES`] (so they fit Durable Object
372/// limits), one partition per batch; one write per batch on a store without
373/// `atomic_multi_key`. A partition's records must be contiguous, as an
374/// export writes them.
375///
376/// On the first record of each partition, [`ImportMode::Fresh`] checks the
377/// partition is empty. After a partition's last record, a store that holds
378/// the `v` class gets `v` = the header's layout version if the export had
379/// no `v` row (a `RefsOnly` export), written only if absent. Records are
380/// plain puts, so an interrupted import can be rerun with
381/// [`ImportMode::Merge`].
382#[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    /// An importer for an export with `header`.
407    ///
408    /// # Errors
409    /// [`StoreError::Unsupported`] if the export's layout version is newer
410    /// than [`keys::LAYOUT_VERSION`], or differs from the implicit layout
411    /// version of a store that cannot hold `v`.
412    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    /// Queue `record`, committing the pending batch first if the record
442    /// changes partition or would overflow it.
443    ///
444    /// # Errors
445    /// [`StoreError::Unsupported`] for a `v` row newer than this binary's
446    /// layout; [`StoreError::Invalid`] for a non-empty partition under
447    /// [`ImportMode::Fresh`]; any error of the store's `apply`.
448    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    /// Commit the current partition's last batch, then its `v` row if it
497    /// needs one.
498    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    /// Commit the pending batches; the number of records imported.
519    pub async fn finish(mut self) -> Result<u64, StoreError> {
520        self.end_partition().await?;
521        Ok(self.imported)
522    }
523}
524
525/// Import a record stream (such as an [`ExportStream`]) taken under
526/// `header`; the number of records imported. Records carry their
527/// partition, so one stream may restore several partitions.
528pub 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
546/// Backend-defined maintenance (R-19). Optional: the pipeline never calls
547/// it. `SQLite` implements it with versioned physical migrations and
548/// `VACUUM INTO` (M0-09); another backend may implement it however it
549/// likes, and every backend still has the portable export above.
550pub trait StoreMaintenance: MaybeSend + MaybeSync {
551    /// The backend's physical layout version (its schema, e.g. `SQLite`'s
552    /// `user_version`); independent of [`keys::LAYOUT_VERSION`].
553    fn layout_version(&self) -> u32;
554
555    /// Migrate the physical layout to the version this binary expects;
556    /// returns the version reached. Idempotent.
557    fn migrate(&self) -> impl Future<Output = Result<u32, StoreError>> + MaybeSend;
558
559    /// Write a consistent backend-native backup to `dest`, a location the
560    /// backend defines (a file path, a bucket URL).
561    fn backup_to(&self, dest: &str) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
562}
563
564/// A future verifiable state root over a partition (R-21). Optional and
565/// unimplemented in the epic; the pipeline never requires it.
566pub trait StateCommitment: NamespaceStore {
567    /// The partition's current root, with the name of its commitment
568    /// scheme; `None` if the backend keeps no root for it.
569    fn root(
570        &self,
571        p: &Partition,
572    ) -> impl Future<Output = Result<Option<(String, Hash)>, StoreError>> + MaybeSend;
573
574    /// A proof of `key`'s value (or absence) against [`Self::root`], in the
575    /// scheme's encoding; `None` if the backend cannot prove it.
576    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    /// Every row of `p`, in pages of 2.
606    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    /// `export_partition` of `p`, encoded.
620    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        // The stream form imports directly, several partitions at once.
684        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        // A `v` row newer than the header claims is refused too.
772        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, &current, 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    /// A store that records the size of every batch it applies.
790    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            // A full store also gets one `v` batch per partition: the
881            // export had no `v` row.
882            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                // Split by bytes (3 big values per batch), by ops and by
886                // partition, never needlessly.
887                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        // Into a full store it gains the header's layout version as `v`.
917        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}