Skip to main content

mkit_server/store/
restore.rs

1//! Restore portable partition exports into a quiesced, fresh store.
2//!
3//! The input is validated and every supplied target partition is checked
4//! before the first write. A caller must provide a newly empty store and keep
5//! traffic off it until restore succeeds: this trait cannot enumerate
6//! unsupplied partitions, and restore is not atomic across partitions. For
7//! relay sources, a successful restore raises each queued sequence and `os`
8//! by `max(snapshot os, supplied target rh)`. Queued sequences then exceed
9//! supplied watermarks. If `os` was absent, the first future allocation is
10//! one above the restored watermark. Fresh targets contain only supplied
11//! partitions; missing relay sources and coordinators are rejected unless
12//! the operator explicitly requests their safe reconstruction.
13//! Already delivered rows no longer exist on the source. Re-keying cannot
14//! fill an older target's missing membership or index rows; R-116 requires
15//! post-restore reconciliation before GA.
16//! Restored coordinator `ls` maxima reset to zero; every restored ref shard
17//! with queued relay rows gets an expired, sweepable `ls` row. The recovery
18//! marker fences watermark reads until R-116 reconciliation completes.
19
20use std::collections::{BTreeMap, BTreeSet};
21
22use super::codec::{self, LeaseRecovery};
23use super::keys::{self, ParsedKey};
24use super::{
25    Batch, BatchOutcome, ExportHeader, ExportReader, ExportRecord, ImportMode, Importer,
26    NamespaceStore, Partition, StoreError, export_page,
27};
28use crate::repo::NamespaceKey;
29use crate::timers::lease_sweep::lease_reference;
30use crate::timers::registry::kinds;
31
32/// Parameters for a restore into empty partitions.
33#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
34pub struct RestoreOptions {
35    /// Require at least this epoch in restored namespace coordinators.
36    pub epoch_at_least: Option<u64>,
37    /// The restore clock reading used for the lease-table recovery marker.
38    pub recovered_at_ms: u64,
39    /// Reconstruct missing sources and coordinators (requires an epoch floor
40    /// for missing coordinators).
41    pub allow_incomplete: bool,
42}
43
44/// Counts committed by a successful restore.
45#[derive(Debug, Clone, Default, PartialEq, Eq)]
46pub struct RestoreReport {
47    /// Number of distinct imported partitions.
48    pub partitions: usize,
49    /// Number of imported records, including the rewritten epoch and sequence.
50    pub records: u64,
51    /// Missing relay sources reconstructed with a sequence floor.
52    pub missing_sources: Vec<Partition>,
53    /// Missing namespace coordinators reconstructed with a new epoch.
54    pub missing_coordinators: Vec<Partition>,
55}
56
57#[derive(Debug)]
58struct SnapshotInfo {
59    index: usize,
60    partition: Partition,
61    header: ExportHeader,
62    epoch: u64,
63    outbox_sequence: u64,
64    has_outbox_sequence: bool,
65    max_relay_sequence: u64,
66    has_relay: bool,
67    sharding_marker: Option<ExportRecord>,
68    inspection_marker: Option<ExportRecord>,
69    authority_fence: bool,
70    authority_generation: Option<u64>,
71}
72
73#[derive(Debug, Clone, Copy, PartialEq, Eq)]
74enum ShardingMode {
75    Single,
76    D34,
77}
78
79impl ShardingMode {
80    fn coordinator(self, partition: &Partition) -> bool {
81        match self {
82            Self::Single => matches!(partition, Partition::Namespace(_)),
83            Self::D34 => matches!(partition, Partition::Coordinator(_)),
84        }
85    }
86}
87
88fn invalid(message: &'static str) -> StoreError {
89    StoreError::Invalid(message.into())
90}
91
92fn corrupt(message: &'static str) -> StoreError {
93    StoreError::Corrupt(message.into())
94}
95
96fn priority(info: &SnapshotInfo) -> u8 {
97    if info.sharding_marker.is_some() {
98        return 0;
99    }
100    match info.partition {
101        Partition::Namespace(_) => 1,
102        Partition::Coordinator(_) => 2,
103        Partition::Ref { .. } => 3,
104        Partition::RefIndex { .. } | Partition::RepoIndex { .. } => 4,
105        Partition::ContentShard(_) => 5,
106    }
107}
108
109fn inspect(
110    index: usize,
111    bytes: &[u8],
112    seen: &mut BTreeSet<Partition>,
113    watermarks: &mut BTreeMap<Partition, u64>,
114) -> Result<SnapshotInfo, StoreError> {
115    let (header, reader) = ExportReader::new(bytes)?;
116    let mut partition = None;
117    let mut epoch = 0;
118    let mut outbox_sequence = 0;
119    let mut has_outbox_sequence = false;
120    let mut max_relay_sequence = 0;
121    let mut has_relay = false;
122    let mut sharding_marker = None;
123    let mut inspection_marker = None;
124    let mut authority_fence = false;
125    let mut authority_generation = None;
126    for record in reader {
127        let record = record?;
128        match &partition {
129            Some(p) if p != &record.partition => {
130                return Err(invalid("snapshot contains multiple partitions"));
131            }
132            None => partition = Some(record.partition.clone()),
133            _ => {}
134        }
135        if record.key == keys::authority_generation() {
136            authority_generation = Some(codec::decode_u64(&record.value)?);
137            authority_fence = true;
138        } else if record.key == keys::epoch_lease() {
139            authority_generation = codec::decode_epoch_lease(&record.value)?.authority_generation;
140            authority_fence |= authority_generation.is_some();
141        } else if record.key == keys::lease_recovery() {
142            authority_fence |=
143                codec::decode_lease_recovery(&record.value)?.authority_fence == Some(true);
144        }
145        validate_inspection_record(&record)?;
146        match keys::parse(&record.key) {
147            Some(ParsedKey::Ticket(_)) => {
148                let ticket = codec::decode_ticket(&record.value)?;
149                authority_generation = authority_generation.max(ticket.authority_generation);
150                authority_fence |= ticket.authority_generation.is_some();
151            }
152            Some(ParsedKey::GrantEpoch) => epoch = codec::decode_u64(&record.value)?,
153            Some(ParsedKey::OutboxSequence) => {
154                outbox_sequence = codec::decode_u64(&record.value)?;
155                has_outbox_sequence = true;
156            }
157            Some(ParsedKey::Relay(seq)) => {
158                has_relay = true;
159                if seq == 0 {
160                    return Err(corrupt("relay sequence is zero"));
161                }
162                codec::decode_relay(&record.value)?;
163                max_relay_sequence = max_relay_sequence.max(seq);
164            }
165            Some(ParsedKey::LayoutVersion)
166                if codec::decode_u32(&record.value)? > keys::LAYOUT_VERSION =>
167            {
168                return Err(StoreError::Unsupported(
169                    "snapshot layout is newer than this binary".into(),
170                ));
171            }
172            Some(ParsedKey::RelayHighWater(source)) => {
173                let rh = codec::decode_u64(&record.value)?;
174                let prior = watermarks.entry(source).or_default();
175                *prior = (*prior).max(rh);
176            }
177            Some(ParsedKey::ShardingMarker) => sharding_marker = Some(record),
178            Some(ParsedKey::InspectionMarker) => inspection_marker = Some(record),
179            _ => {}
180        }
181    }
182    let partition = partition.ok_or_else(|| invalid("snapshot contains no partition records"))?;
183    if !seen.insert(partition.clone()) {
184        return Err(invalid("duplicate partition snapshot"));
185    }
186    validate_sequences(max_relay_sequence, outbox_sequence, has_outbox_sequence)?;
187    if sharding_marker.is_some()
188        && partition != Partition::Namespace(NamespaceKey::deployment_default())
189    {
190        return Err(corrupt("sharding marker outside root namespace partition"));
191    }
192    if inspection_marker.is_some()
193        && partition != Partition::Namespace(NamespaceKey::deployment_default())
194    {
195        return Err(corrupt(
196            "inspection marker outside root namespace partition",
197        ));
198    }
199    Ok(SnapshotInfo {
200        index,
201        partition,
202        header,
203        epoch,
204        outbox_sequence,
205        has_outbox_sequence,
206        max_relay_sequence,
207        has_relay,
208        sharding_marker,
209        inspection_marker,
210        authority_fence,
211        authority_generation,
212    })
213}
214
215fn validate_sequences(
216    max_relay_sequence: u64,
217    outbox_sequence: u64,
218    has_outbox_sequence: bool,
219) -> Result<(), StoreError> {
220    if max_relay_sequence > outbox_sequence {
221        return Err(corrupt("relay sequence exceeds outbox sequence"));
222    }
223    if max_relay_sequence > 0 && outbox_sequence == 0 {
224        return Err(corrupt("relay rows require an outbox sequence"));
225    }
226    if has_outbox_sequence && outbox_sequence == 0 {
227        return Err(corrupt("outbox sequence is zero"));
228    }
229    Ok(())
230}
231
232fn validate_inspection_record(record: &ExportRecord) -> Result<(), StoreError> {
233    match keys::parse(&record.key) {
234        Some(ParsedKey::InspectionMarker) if record.value.as_bytes() != b"on" => {
235            Err(corrupt("invalid inspection marker"))
236        }
237        Some(ParsedKey::InspectionFlag { repo, id }) => {
238            inspection_partition(&record.partition, &repo, true)?;
239            if super::inspection_flags::decode_flag(&record.value)?.id != id {
240                return Err(corrupt("inspection flag id disagrees with key"));
241            }
242            Ok(())
243        }
244        Some(ParsedKey::InspectionVersion(repo)) => {
245            inspection_partition(&record.partition, &repo, true)?;
246            if codec::decode_u64(&record.value)? == 0 {
247                return Err(corrupt("inspection registry version is zero"));
248            }
249            Ok(())
250        }
251        Some(ParsedKey::InspectionHold { repo, .. }) => {
252            inspection_partition(&record.partition, &repo, true)?;
253            if !record.value.as_bytes().is_empty() {
254                return Err(corrupt("invalid inspection hold row"));
255            }
256            Ok(())
257        }
258        Some(ParsedKey::InspectionHoldIndex { repo, .. }) => {
259            inspection_partition(&record.partition, &repo, false)?;
260            super::inspection_holds::validate_advance_hold(&record.value)
261        }
262        Some(ParsedKey::InspectionHoldManifest { repo, .. }) => {
263            inspection_partition(&record.partition, &repo, true)?;
264            super::inspection_holds::validate_manifest(&record.value)
265        }
266        _ => Ok(()),
267    }
268}
269
270fn inspection_partition(
271    partition: &Partition,
272    repo: &crate::RepoName,
273    registry: bool,
274) -> Result<(), StoreError> {
275    match partition {
276        Partition::Namespace(_) => Ok(()),
277        Partition::RepoIndex {
278            repo: stored,
279            prefix: 0,
280            ..
281        } if registry && stored == repo => Ok(()),
282        Partition::Ref { repo: stored, .. } if !registry && stored == repo => Ok(()),
283        _ => Err(corrupt("inspection row outside its authority partition")),
284    }
285}
286
287fn should_drop(record: &ExportRecord) -> bool {
288    record.key == keys::revoke_cursor(false)
289        || record.key == keys::revoke_cursor(true)
290        || record.key == keys::backup_state()
291        || record.key == keys::epoch_lease()
292        || record.key == keys::relay_scan()
293        || record.key == keys::lease_reconcile()
294        || matches!(
295            keys::parse(&record.key),
296            Some(ParsedKey::Timer { kind, .. }) if kind == kinds::BACKUP.get() || kind == kinds::LEASE_SWEEP.get()
297        )
298}
299
300/// Write `lr 00` exactly as the pipeline's recovery declaration does.
301///
302/// This must finish before a restored coordinator serves writes.
303pub async fn mark_lease_table_recovered<S: NamespaceStore>(
304    store: &S,
305    partition: &Partition,
306    recovered_at_ms: u64,
307) -> Result<(), StoreError> {
308    let key = keys::lease_recovery();
309    let prior = store.get(partition, &key).await?;
310    let authority = store.get(partition, &keys::authority_generation()).await?;
311    if let Some(value) = authority.as_ref() {
312        codec::decode_u64(value)?;
313    }
314    let mode = prior
315        .as_ref()
316        .map(codec::decode_lease_recovery)
317        .transpose()?;
318    if mode.is_some_and(|m| m.authority_fence == Some(true)) && authority.is_none() {
319        return Err(corrupt("fenced recovery requires authority generation"));
320    }
321    let batch = Batch::new()
322        .require(match authority.as_ref() {
323            Some(value) => super::Precondition::Equals(keys::authority_generation(), value.clone()),
324            None => super::Precondition::Absent(keys::authority_generation()),
325        })
326        .require(match prior.as_ref() {
327            Some(value) => super::Precondition::Equals(key.clone(), value.clone()),
328            None => super::Precondition::Absent(key.clone()),
329        })
330        .delete(keys::lease_reconcile())
331        .delete(keys::revoke_cursor(false))
332        .delete(keys::revoke_cursor(true))
333        .put(
334            key,
335            codec::encode_lease_recovery(&LeaseRecovery {
336                authority_fence: authority.as_ref().map(|_| true),
337                authority_ready: authority
338                    .as_ref()
339                    .map(|_| mode.is_some_and(|m| m.authority_ready == Some(true))),
340                activation_only: None,
341                resumed_at_ms: recovered_at_ms,
342            }),
343        );
344    match store.apply(partition, batch).await? {
345        BatchOutcome::Committed => Ok(()),
346        _ => Err(StoreError::unavailable(
347            "recovery marker batch did not commit",
348        )),
349    }
350}
351
352struct RestorePlan {
353    infos: Vec<SnapshotInfo>,
354    shifts: BTreeMap<Partition, u64>,
355    mode: ShardingMode,
356    missing_sources: Vec<(Partition, u64)>,
357    missing_coordinators: Vec<Partition>,
358}
359
360const EPOCH_RESTORE_JUMP: u64 = 1 << 32;
361
362fn namespace(partition: &Partition) -> Option<&NamespaceKey> {
363    match partition {
364        Partition::Namespace(ns)
365        | Partition::Coordinator(ns)
366        | Partition::Ref { ns, .. }
367        | Partition::RefIndex { ns, .. }
368        | Partition::RepoIndex { ns, .. } => Some(ns),
369        Partition::ContentShard(_) => None,
370    }
371}
372
373struct Completeness {
374    sources: BTreeSet<Partition>,
375    missing_sources: Vec<(Partition, u64)>,
376    missing_coordinators: Vec<Partition>,
377}
378
379fn check_completeness(
380    infos: &[SnapshotInfo],
381    seen: &BTreeSet<Partition>,
382    watermarks: &BTreeMap<Partition, u64>,
383    mode: ShardingMode,
384    opts: RestoreOptions,
385) -> Result<Completeness, StoreError> {
386    let sources: BTreeSet<_> = infos
387        .iter()
388        .filter(|info| info.has_outbox_sequence || info.has_relay)
389        .map(|info| info.partition.clone())
390        .collect();
391    let missing_sources: Vec<_> = watermarks
392        .iter()
393        .filter(|(source, _)| !sources.contains(*source))
394        .map(|(source, floor)| (source.clone(), *floor))
395        .collect();
396    let required_coordinators: BTreeSet<_> = infos
397        .iter()
398        .map(|info| &info.partition)
399        .chain(missing_sources.iter().map(|(source, _)| source))
400        .filter(|partition| match mode {
401            ShardingMode::Single => matches!(partition, Partition::Namespace(_)),
402            ShardingMode::D34 => matches!(
403                partition,
404                Partition::Ref { .. } | Partition::RefIndex { .. } | Partition::RepoIndex { .. }
405            ),
406        })
407        .filter_map(namespace)
408        .map(|ns| match mode {
409            ShardingMode::Single => Partition::Namespace(ns.clone()),
410            ShardingMode::D34 => Partition::Coordinator(ns.clone()),
411        })
412        .collect();
413    let missing_coordinators: Vec<_> = required_coordinators.difference(seen).cloned().collect();
414    if !missing_sources.is_empty() && !opts.allow_incomplete {
415        return Err(StoreError::Invalid(
416            format!("missing relay sources: {missing_sources:?}").into(),
417        ));
418    }
419    if !missing_coordinators.is_empty() && (!opts.allow_incomplete || opts.epoch_at_least.is_none())
420    {
421        return Err(StoreError::Invalid(format!("missing namespace coordinators (requires --allow-incomplete and --epoch-at-least): {missing_coordinators:?}").into()));
422    }
423    // A missing coordinator's history is unknown, so the floor must exceed any
424    // epoch a namespace plausibly issued; grants compare epochs by equality.
425    if !missing_coordinators.is_empty()
426        && opts
427            .epoch_at_least
428            .is_some_and(|floor| floor < EPOCH_RESTORE_JUMP)
429    {
430        return Err(StoreError::Invalid(
431            format!(
432                "--epoch-at-least must be at least {EPOCH_RESTORE_JUMP} for missing coordinators"
433            )
434            .into(),
435        ));
436    }
437    for partition in sources
438        .iter()
439        .chain(missing_sources.iter().map(|(partition, _)| partition))
440    {
441        if !matches!(partition, Partition::Ref { .. }) {
442            return Err(invalid("relay high-water names a non-ref source"));
443        }
444    }
445    Ok(Completeness {
446        sources,
447        missing_sources,
448        missing_coordinators,
449    })
450}
451
452#[allow(clippy::too_many_lines)] // All portable input, fence and arithmetic checks must finish before any import writes.
453async fn prepare<S: NamespaceStore>(
454    snapshots: &[Vec<u8>],
455    target: &S,
456    opts: RestoreOptions,
457) -> Result<RestorePlan, StoreError> {
458    let mut seen = BTreeSet::new();
459    let mut watermarks = BTreeMap::new();
460    let mut infos = Vec::new();
461    for (index, bytes) in snapshots.iter().enumerate() {
462        let (_, mut reader) = ExportReader::new(bytes)?;
463        if reader.next().transpose()?.is_none() {
464            tracing::warn!(index, "skipping recordless restore archive");
465            continue;
466        }
467        infos.push(inspect(index, bytes, &mut seen, &mut watermarks)?);
468    }
469    let marker = infos
470        .iter()
471        .find_map(|info| info.sharding_marker.as_ref())
472        .ok_or_else(|| invalid("restore requires the root sharding marker"))?;
473    let mode = match marker.value.as_bytes() {
474        b"single" => ShardingMode::Single,
475        b"d34" => ShardingMode::D34,
476        _ => return Err(corrupt("invalid sharding marker")),
477    };
478    let Completeness {
479        sources,
480        missing_sources,
481        missing_coordinators,
482    } = check_completeness(&infos, &seen, &watermarks, mode, opts)?;
483    for info in &infos {
484        let incompatible = match mode {
485            ShardingMode::Single => matches!(
486                info.partition,
487                Partition::Coordinator(_)
488                    | Partition::Ref { .. }
489                    | Partition::RepoIndex { .. }
490                    | Partition::RefIndex { .. }
491            ),
492            ShardingMode::D34 => {
493                matches!(info.partition, Partition::Namespace(_)) && info.sharding_marker.is_none()
494            }
495        };
496        if incompatible {
497            return Err(invalid("snapshot partition conflicts with sharding marker"));
498        }
499    }
500    for info in infos.iter().filter(|info| info.authority_fence) {
501        let owner = infos.iter().find(|candidate| {
502            mode.coordinator(&candidate.partition)
503                && namespace(&candidate.partition) == namespace(&info.partition)
504        });
505        if owner.is_none_or(|owner| {
506            owner.authority_generation.is_none()
507                || info.authority_generation > owner.authority_generation
508        }) {
509            return Err(invalid(
510                "fenced snapshot requires authoritative generation without rollback",
511            ));
512        }
513    }
514    for coordinator in &missing_coordinators {
515        if infos.iter().any(|info| {
516            info.authority_fence && namespace(&info.partition) == namespace(coordinator)
517        }) {
518            return Err(invalid(
519                "fenced namespace coordinator cannot be reconstructed without authority state",
520            ));
521        }
522    }
523    for partition in missing_sources
524        .iter()
525        .map(|(p, _)| p)
526        .chain(missing_coordinators.iter())
527    {
528        let page = export_page(target, partition, None, 1).await?;
529        if !page.records.is_empty() {
530            return Err(invalid("restore target partition is not empty"));
531        }
532    }
533
534    // Preflight every supplied partition before the first Importer batch.
535    for info in &infos {
536        Importer::new(target, &info.header, ImportMode::Fresh)?;
537        let page = export_page(target, &info.partition, None, 1).await?;
538        if !page.records.is_empty() {
539            return Err(invalid("restore target partition is not empty"));
540        }
541    }
542
543    // Check all arithmetic before writing anything. Every shifted row and
544    // future allocation then starts above the largest supplied target rh.
545    let mut shifts = BTreeMap::new();
546    for info in &infos {
547        if sources.contains(&info.partition) {
548            let shift = info
549                .outbox_sequence
550                .max(*watermarks.get(&info.partition).unwrap_or(&0));
551            info.outbox_sequence
552                .checked_add(shift)
553                .and_then(|next| next.checked_add(1))
554                .ok_or_else(|| invalid("restore relay sequence overflow"))?;
555            info.max_relay_sequence
556                .checked_add(shift)
557                .ok_or_else(|| invalid("restore relay sequence overflow"))?;
558            shifts.insert(info.partition.clone(), shift);
559        }
560        if mode.coordinator(&info.partition) {
561            info.epoch
562                .checked_add(EPOCH_RESTORE_JUMP)
563                .ok_or_else(|| invalid("restore epoch overflow"))?;
564        }
565    }
566
567    infos.sort_by_key(|info| (priority(info), info.partition.clone()));
568    Ok(RestorePlan {
569        infos,
570        shifts,
571        mode,
572        missing_sources,
573        missing_coordinators,
574    })
575}
576
577async fn import_one<S: NamespaceStore>(
578    bytes: &[u8],
579    target: &S,
580    info: &SnapshotInfo,
581    shift: Option<u64>,
582    mode: ShardingMode,
583    opts: RestoreOptions,
584) -> Result<u64, StoreError> {
585    let (_, reader) = ExportReader::new(bytes)?;
586    let mut importer = Importer::new(target, &info.header, ImportMode::Fresh)?;
587    if let Some(marker) = &info.sharding_marker {
588        importer.push(marker.clone()).await?;
589    }
590    if let Some(marker) = &info.inspection_marker {
591        importer.push(marker.clone()).await?;
592    }
593    for record in reader {
594        let mut record = record?;
595        if record.key == keys::sharding_marker() && info.sharding_marker.is_some() {
596            continue;
597        }
598        if record.key == keys::inspection_marker() && info.inspection_marker.is_some() {
599            continue;
600        }
601        if should_drop(&record)
602            || (mode.coordinator(&info.partition) && record.key == keys::grant_epoch())
603        {
604            continue;
605        }
606        if let Some(shift) = shift {
607            match keys::parse(&record.key) {
608                Some(ParsedKey::Relay(seq)) => {
609                    record.key = keys::relay(
610                        seq.checked_add(shift)
611                            .ok_or_else(|| invalid("restore relay sequence overflow"))?,
612                    );
613                }
614                Some(ParsedKey::OutboxSequence) => {
615                    record.value = codec::encode_u64(
616                        info.outbox_sequence
617                            .checked_add(shift)
618                            .ok_or_else(|| invalid("restore relay sequence overflow"))?,
619                    );
620                }
621                _ => {}
622            }
623        }
624        if let Some(ParsedKey::LeasedShard { repo, shard_ref }) = keys::parse(&record.key) {
625            let mut row = codec::decode_leased_shard(&record.value)?;
626            row.expires_at_ms = opts.recovered_at_ms;
627            row.relay_watermark_ms = 0;
628            row.sweep_due_ms = opts.recovered_at_ms;
629            record.value = codec::encode_leased_shard(&row);
630            importer.push(record).await?;
631            importer
632                .push(ExportRecord::new(
633                    info.partition.clone(),
634                    keys::timer(
635                        opts.recovered_at_ms,
636                        kinds::LEASE_SWEEP.get(),
637                        &lease_reference(&repo, &shard_ref),
638                    ),
639                    super::Value::default(),
640                ))
641                .await?;
642            continue;
643        }
644        importer.push(record).await?;
645    }
646    if mode.coordinator(&info.partition) {
647        let epoch = info
648            .epoch
649            .checked_add(EPOCH_RESTORE_JUMP)
650            .ok_or_else(|| invalid("restore epoch overflow"))?
651            .max(opts.epoch_at_least.unwrap_or(0));
652        importer
653            .push(ExportRecord::new(
654                info.partition.clone(),
655                keys::grant_epoch(),
656                codec::encode_u64(epoch),
657            ))
658            .await?;
659    }
660    if let Some(shift) = shift
661        && !info.has_outbox_sequence
662        && shift > 0
663    {
664        importer
665            .push(ExportRecord::new(
666                info.partition.clone(),
667                keys::outbox_sequence(),
668                codec::encode_u64(shift),
669            ))
670            .await?;
671    }
672    let records = importer.finish().await?;
673    if mode.coordinator(&info.partition) {
674        mark_lease_table_recovered(target, &info.partition, opts.recovered_at_ms).await?;
675    }
676    Ok(records)
677}
678
679async fn create_missing_sources<S: NamespaceStore>(
680    target: &S,
681    sources: &[(Partition, u64)],
682    supplied: &BTreeSet<Partition>,
683    report: &mut RestoreReport,
684) -> Result<(), StoreError> {
685    for (partition, floor) in sources {
686        if *floor > 0 {
687            match target
688                .apply(
689                    partition,
690                    Batch::new().put(keys::outbox_sequence(), codec::encode_u64(*floor)),
691                )
692                .await?
693            {
694                BatchOutcome::Committed => {}
695                _ => {
696                    return Err(StoreError::unavailable(
697                        "missing relay source creation did not commit",
698                    ));
699                }
700            }
701        }
702        report.records += u64::from(*floor > 0);
703        if !supplied.contains(partition) {
704            report.partitions += 1;
705        }
706        report.missing_sources.push(partition.clone());
707    }
708    Ok(())
709}
710
711async fn seed_relay_leases<S: NamespaceStore>(
712    target: &S,
713    sources: Vec<Partition>,
714    recovered_at_ms: u64,
715    report: &mut RestoreReport,
716) -> Result<(), StoreError> {
717    for source in sources {
718        let Partition::Ref {
719            ns,
720            repo,
721            shard_ref,
722        } = source
723        else {
724            unreachable!("relay sources were filtered to ref shards")
725        };
726        let coordinator = Partition::Coordinator(ns);
727        let key = keys::leased_shard(&repo, &shard_ref);
728        if target.get(&coordinator, &key).await?.is_some() {
729            continue;
730        }
731        let epoch = target
732            .get(&coordinator, &keys::grant_epoch())
733            .await?
734            .as_ref()
735            .map(codec::decode_u64)
736            .transpose()?
737            .unwrap_or(0);
738        let authority_generation = target
739            .get(&coordinator, &keys::authority_generation())
740            .await?
741            .as_ref()
742            .map(codec::decode_u64)
743            .transpose()?;
744        let row = codec::LeasedShard {
745            authority_generation,
746            acked_authority_generation: authority_generation,
747            epoch,
748            expires_at_ms: recovered_at_ms,
749            acked_epoch: epoch,
750            relay_watermark_ms: 0,
751            sweep_due_ms: recovered_at_ms,
752        };
753        let batch = Batch::new().put(key, codec::encode_leased_shard(&row)).put(
754            keys::timer(
755                recovered_at_ms,
756                kinds::LEASE_SWEEP.get(),
757                &lease_reference(&repo, &shard_ref),
758            ),
759            super::Value::default(),
760        );
761        match target.apply(&coordinator, batch).await? {
762            BatchOutcome::Committed => report.records += 2,
763            _ => {
764                return Err(StoreError::Corrupt(
765                    "restored relay lease row did not commit".into(),
766                ));
767            }
768        }
769    }
770    Ok(())
771}
772
773/// Restore exports into a fresh store, in dependency order.
774///
775/// Each nonempty export must contain exactly one partition. Recordless
776/// exports are skipped with a warning because their header cannot identify
777/// a target partition. All input and all supplied target
778/// emptiness checks complete before any partition is written. This function
779/// is intended for an offline, newly empty target: imports span batches.
780///
781/// # Errors
782/// Malformed or duplicate snapshots, non-empty supplied target partitions,
783/// epoch or relay sequence overflow, or a store failure.
784pub async fn restore<S: NamespaceStore>(
785    snapshots: &[Vec<u8>],
786    target: &S,
787    opts: RestoreOptions,
788) -> Result<RestoreReport, StoreError> {
789    let plan = prepare(snapshots, target, opts).await?;
790    let supplied: BTreeSet<_> = plan
791        .infos
792        .iter()
793        .map(|info| info.partition.clone())
794        .collect();
795    let relay_sources: Vec<Partition> = plan
796        .infos
797        .iter()
798        .filter(|info| info.has_relay && matches!(info.partition, Partition::Ref { .. }))
799        .map(|info| info.partition.clone())
800        .collect();
801    let mut report = RestoreReport::default();
802    let mut wrote_missing_coordinators = false;
803    let mut wrote_missing_sources = false;
804    for info in plan.infos {
805        if !wrote_missing_coordinators && priority(&info) > 2 {
806            for partition in &plan.missing_coordinators {
807                let epoch = opts
808                    .epoch_at_least
809                    .ok_or_else(|| invalid("missing coordinator epoch floor"))?;
810                match target
811                    .apply(
812                        partition,
813                        Batch::new().put(keys::grant_epoch(), codec::encode_u64(epoch)),
814                    )
815                    .await?
816                {
817                    BatchOutcome::Committed => {}
818                    _ => {
819                        return Err(StoreError::unavailable(
820                            "missing coordinator creation did not commit",
821                        ));
822                    }
823                }
824                mark_lease_table_recovered(target, partition, opts.recovered_at_ms).await?;
825                report.records += 2;
826                report.partitions += 1;
827                report.missing_coordinators.push(partition.clone());
828            }
829            wrote_missing_coordinators = true;
830        }
831        if !wrote_missing_sources && priority(&info) > 3 {
832            create_missing_sources(target, &plan.missing_sources, &supplied, &mut report).await?;
833            wrote_missing_sources = true;
834        }
835        report.records += import_one(
836            &snapshots[info.index],
837            target,
838            &info,
839            plan.shifts.get(&info.partition).copied(),
840            plan.mode,
841            opts,
842        )
843        .await?;
844        report.partitions += 1;
845    }
846    if !wrote_missing_sources {
847        create_missing_sources(target, &plan.missing_sources, &supplied, &mut report).await?;
848    }
849    seed_relay_leases(target, relay_sources, opts.recovered_at_ms, &mut report).await?;
850    Ok(report)
851}
852
853#[cfg(test)]
854mod tests {
855    use std::sync::Mutex;
856
857    use futures_executor::block_on;
858
859    use super::*;
860    use crate::memory::MemoryKv;
861    use crate::relay::{NoHook, RelayBudget, RelayHandler};
862    use crate::repo::{NamespaceKey, RepoName};
863    use crate::store::{
864        Cursor, EXPORT_END, Key, PartitionStats, ScanPage, StoreCapabilities, Value, Write,
865        encode_export_header, encode_export_record,
866    };
867    use crate::timers::{DueTimer, TimerCtx, TimerHandler, registry::kinds};
868
869    fn root() -> Partition {
870        Partition::Namespace(NamespaceKey::deployment_default())
871    }
872
873    fn coordinator() -> Partition {
874        Partition::Coordinator(NamespaceKey::deployment_default())
875    }
876
877    fn source() -> Partition {
878        Partition::Ref {
879            ns: NamespaceKey::deployment_default(),
880            repo: RepoName::new("repo").unwrap(),
881            shard_ref: "refs/heads/main".into(),
882        }
883    }
884
885    fn target() -> Partition {
886        Partition::RepoIndex {
887            ns: NamespaceKey::deployment_default(),
888            repo: RepoName::new("repo").unwrap(),
889            prefix: 1,
890        }
891    }
892
893    fn empty_target() -> Partition {
894        Partition::ContentShard(7)
895    }
896
897    fn old_lease() -> Value {
898        codec::encode_epoch_lease(&codec::EpochLease {
899            authority_ready: None,
900            authority_generation: None,
901            epoch: 7,
902            expires_at_ms: 999_999,
903            config_version: 1,
904        })
905    }
906
907    fn snapshot(partition: &Partition, rows: Vec<(Key, Value)>) -> Vec<u8> {
908        let header = ExportHeader::new(keys::LAYOUT_VERSION, 100);
909        let mut bytes = encode_export_header(&header).to_vec();
910        for (key, value) in rows {
911            bytes.extend_from_slice(
912                &encode_export_record(&ExportRecord::new(partition.clone(), key, value)).unwrap(),
913            );
914        }
915        bytes.extend_from_slice(&EXPORT_END);
916        bytes
917    }
918
919    #[derive(Default)]
920    struct RootFirst {
921        inner: MemoryKv,
922        writes: Mutex<Vec<Partition>>,
923    }
924
925    impl NamespaceStore for RootFirst {
926        fn capabilities(&self) -> StoreCapabilities {
927            self.inner.capabilities()
928        }
929
930        async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
931            self.inner.get(p, key).await
932        }
933
934        async fn scan(
935            &self,
936            p: &Partition,
937            start: &Key,
938            end: &Key,
939            after: Option<&Cursor>,
940            limit: u32,
941        ) -> Result<ScanPage, StoreError> {
942            self.inner.scan(p, start, end, after, limit).await
943        }
944
945        async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
946            if p == &root() && self.writes.lock().unwrap().is_empty() {
947                assert!(matches!(
948                    batch.writes.first(),
949                    Some(Write::Put(key, _)) if key == &keys::sharding_marker()
950                ));
951            }
952            if p != &root()
953                && self
954                    .inner
955                    .get(&root(), &keys::sharding_marker())
956                    .await?
957                    .is_none()
958            {
959                return Err(invalid("root marker was not imported first"));
960            }
961            self.writes.lock().unwrap().push(p.clone());
962            self.inner.apply(p, batch).await
963        }
964
965        async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
966            self.inner.stats(p).await
967        }
968
969        async fn probe(&self) -> Result<(), StoreError> {
970            self.inner.probe().await
971        }
972    }
973
974    fn assert_delivers_once(store: RootFirst, shift: u64) {
975        let handler = RelayHandler {
976            target: store,
977            hook: NoHook,
978            budget: RelayBudget::default(),
979        };
980        let source_partition = source();
981        let ctx = TimerCtx {
982            store: &handler.target,
983            partition: &source_partition,
984            now_ms: 200,
985        };
986        let timer = DueTimer {
987            due_at_ms: 200,
988            kind: kinds::RELAY,
989            reference: bytes::Bytes::default(),
990            value: Value::default(),
991        };
992        block_on(handler.fire(&ctx, &timer)).unwrap();
993        for (partition, key, seq) in [
994            (target(), Key::new(b"m\0x".to_vec()), 2 + shift),
995            (empty_target(), Key::new(b"m\0y".to_vec()), 3 + shift),
996        ] {
997            assert_eq!(
998                block_on(handler.target.get(&partition, &key)).unwrap(),
999                Some(Value::default())
1000            );
1001            assert_eq!(
1002                block_on(
1003                    handler
1004                        .target
1005                        .get(&partition, &keys::relay_high_water(&source()).unwrap())
1006                )
1007                .unwrap(),
1008                Some(codec::encode_u64(seq))
1009            );
1010        }
1011        let count_target_applies = || {
1012            handler
1013                .target
1014                .writes
1015                .lock()
1016                .unwrap()
1017                .iter()
1018                .filter(|p| **p == target() || **p == empty_target())
1019                .count()
1020        };
1021        let first = count_target_applies();
1022        block_on(handler.fire(&ctx, &timer)).unwrap();
1023        assert_eq!(
1024            count_target_applies(),
1025            first,
1026            "repeat fire has no second target effect"
1027        );
1028    }
1029
1030    fn assert_restored_lease_rows(store: &RootFirst) {
1031        for name in ["other", "repo"] {
1032            let repo = RepoName::new(name).expect("test repository name");
1033            let row = block_on(store.get(
1034                &coordinator(),
1035                &keys::leased_shard(&repo, "refs/heads/main"),
1036            ))
1037            .expect("read restored lease")
1038            .expect("restored lease exists");
1039            let row = codec::decode_leased_shard(&row).expect("decode restored lease");
1040            assert_eq!(row.relay_watermark_ms, 0);
1041            assert_eq!(row.expires_at_ms, 500);
1042            assert_eq!(row.sweep_due_ms, 500);
1043            assert!(
1044                block_on(store.get(
1045                    &coordinator(),
1046                    &keys::timer(
1047                        500,
1048                        kinds::LEASE_SWEEP.get(),
1049                        &lease_reference(&repo, "refs/heads/main")
1050                    )
1051                ))
1052                .expect("read restored sweep timer")
1053                .is_some()
1054            );
1055        }
1056        assert!(
1057            block_on(store.get(
1058                &coordinator(),
1059                &keys::timer(
1060                    999_999,
1061                    kinds::LEASE_SWEEP.get(),
1062                    &lease_reference(
1063                        &RepoName::new("other").expect("test repository name"),
1064                        "refs/heads/main"
1065                    )
1066                )
1067            ))
1068            .expect("read old sweep timer")
1069            .is_none()
1070        );
1071    }
1072
1073    fn assert_restored_source_rows(
1074        store: RootFirst,
1075        relay: Value,
1076        relay_to_empty: Value,
1077        shift: u64,
1078    ) {
1079        assert_eq!(
1080            block_on(store.get(&source(), &keys::outbox_sequence()))
1081                .expect("read restored source sequence"),
1082            Some(codec::encode_u64(3 + shift))
1083        );
1084        assert_eq!(
1085            block_on(store.get(&source(), &keys::relay(2 + shift)))
1086                .expect("read first restored relay"),
1087            Some(relay)
1088        );
1089        assert_eq!(
1090            block_on(store.get(&source(), &keys::relay(3 + shift)))
1091                .expect("read second restored relay"),
1092            Some(relay_to_empty)
1093        );
1094        for key in [
1095            keys::relay(2),
1096            keys::epoch_lease(),
1097            keys::relay_scan(),
1098            keys::backup_state(),
1099            keys::timer(123, kinds::BACKUP.get(), b""),
1100        ] {
1101            assert_eq!(
1102                block_on(store.get(&source(), &key)).expect("read discarded source row"),
1103                None
1104            );
1105        }
1106        assert_delivers_once(store, shift);
1107    }
1108
1109    #[test]
1110    fn fresh_restore_orders_root_and_rewrites_recovery_epoch_and_relay() {
1111        let index = snapshot(
1112            &target(),
1113            vec![(
1114                keys::relay_high_water(&source()).unwrap(),
1115                codec::encode_u64(10),
1116            )],
1117        );
1118        let relay_to = |target, key: &[u8]| {
1119            codec::encode_relay(&codec::RelayV1 {
1120                at_ms: 100,
1121                target,
1122                puts: vec![(Key::new(key.to_vec()), Value::default())],
1123                deletes: Vec::new(),
1124            })
1125            .unwrap()
1126        };
1127        let relay = relay_to(target(), b"m\0x");
1128        let relay_to_empty = relay_to(empty_target(), b"m\0y");
1129        let ref_shard = snapshot(
1130            &source(),
1131            vec![
1132                (keys::outbox_sequence(), codec::encode_u64(3)),
1133                (keys::relay(2), relay.clone()),
1134                (keys::relay(3), relay_to_empty.clone()),
1135                (keys::epoch_lease(), old_lease()),
1136                (keys::relay_scan(), Value::default()),
1137                (keys::backup_state(), Value::default()),
1138                (keys::timer(123, kinds::BACKUP.get(), b""), Value::default()),
1139            ],
1140        );
1141        let coordinator_snapshot = snapshot(
1142            &coordinator(),
1143            vec![
1144                (keys::grant_epoch(), codec::encode_u64(7)),
1145                (
1146                    keys::leased_shard(&RepoName::new("other").unwrap(), "refs/heads/main"),
1147                    codec::encode_leased_shard(&codec::LeasedShard {
1148                        authority_generation: None,
1149                        acked_authority_generation: None,
1150                        epoch: 7,
1151                        expires_at_ms: 999_999,
1152                        acked_epoch: 7,
1153                        relay_watermark_ms: 888_888,
1154                        sweep_due_ms: 999_999,
1155                    }),
1156                ),
1157                (
1158                    keys::timer(
1159                        999_999,
1160                        kinds::LEASE_SWEEP.get(),
1161                        &lease_reference(&RepoName::new("other").unwrap(), "refs/heads/main"),
1162                    ),
1163                    Value::default(),
1164                ),
1165            ],
1166        );
1167        let root_snapshot = snapshot(
1168            &root(),
1169            vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1170        );
1171        let store = RootFirst::default();
1172        let report = block_on(restore(
1173            &[index, ref_shard, coordinator_snapshot, root_snapshot],
1174            &store,
1175            RestoreOptions {
1176                epoch_at_least: Some(11),
1177                recovered_at_ms: 500,
1178                allow_incomplete: false,
1179            },
1180        ))
1181        .unwrap();
1182        assert_eq!(report.partitions, 4);
1183        assert_eq!(store.writes.lock().unwrap().first(), Some(&root()));
1184        assert_eq!(
1185            block_on(store.get(&coordinator(), &keys::grant_epoch())).unwrap(),
1186            Some(codec::encode_u64(7 + EPOCH_RESTORE_JUMP))
1187        );
1188        let lr = block_on(store.get(&coordinator(), &keys::lease_recovery()))
1189            .unwrap()
1190            .unwrap();
1191        assert_eq!(
1192            codec::decode_lease_recovery(&lr).unwrap().resumed_at_ms,
1193            500
1194        );
1195        assert_restored_lease_rows(&store);
1196        let writes = store.writes.lock().unwrap();
1197        let coordinator_write = writes.iter().position(|p| p == &coordinator()).unwrap();
1198        let ref_write = writes.iter().position(|p| p == &source()).unwrap();
1199        assert!(coordinator_write < ref_write);
1200        drop(writes);
1201        assert_eq!(
1202            block_on(store.get(&root(), &keys::grant_epoch())).unwrap(),
1203            None,
1204            "the D34 root marker is not a coordinator"
1205        );
1206        assert_eq!(
1207            block_on(store.get(&root(), &keys::lease_recovery())).unwrap(),
1208            None
1209        );
1210        assert_restored_source_rows(store, relay, relay_to_empty, 10);
1211    }
1212
1213    #[test]
1214    fn fresh_preflight_refuses_nonempty_without_partial_root_import() {
1215        let root_snapshot = snapshot(
1216            &root(),
1217            vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1218        );
1219        let coordinator_snapshot = snapshot(
1220            &coordinator(),
1221            vec![(keys::grant_epoch(), codec::encode_u64(8))],
1222        );
1223        let store = MemoryKv::default();
1224        block_on(store.apply(
1225            &coordinator(),
1226            Batch::new().put(keys::grant_epoch(), codec::encode_u64(9)),
1227        ))
1228        .unwrap();
1229        assert!(matches!(
1230            block_on(restore(
1231                &[root_snapshot, coordinator_snapshot],
1232                &store,
1233                RestoreOptions::default()
1234            )),
1235            Err(StoreError::Invalid(_))
1236        ));
1237        assert_eq!(
1238            block_on(store.get(&root(), &keys::sharding_marker())).unwrap(),
1239            None
1240        );
1241    }
1242
1243    #[test]
1244    fn malformed_and_overflowing_snapshots_fail_before_writes() {
1245        let root_snapshot = snapshot(
1246            &root(),
1247            vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1248        );
1249        let store = MemoryKv::default();
1250        assert!(matches!(
1251            block_on(restore(&[], &store, RestoreOptions::default())),
1252            Err(StoreError::Invalid(_))
1253        ));
1254        assert!(matches!(
1255            block_on(restore(
1256                &[root_snapshot.clone(), root_snapshot.clone()],
1257                &store,
1258                RestoreOptions::default()
1259            )),
1260            Err(StoreError::Invalid(_))
1261        ));
1262        let empty = [
1263            encode_export_header(&ExportHeader::new(keys::LAYOUT_VERSION, 0)).as_ref(),
1264            &EXPORT_END,
1265        ]
1266        .concat();
1267        assert!(matches!(
1268            block_on(restore(&[empty], &store, RestoreOptions::default())),
1269            Err(StoreError::Invalid(_))
1270        ));
1271        let overflowing = snapshot(
1272            &source(),
1273            vec![(keys::outbox_sequence(), codec::encode_u64(u64::MAX))],
1274        );
1275        let coordinator_snapshot = snapshot(
1276            &coordinator(),
1277            vec![(keys::grant_epoch(), codec::encode_u64(1))],
1278        );
1279        let err = block_on(restore(
1280            &[root_snapshot, coordinator_snapshot, overflowing],
1281            &store,
1282            RestoreOptions::default(),
1283        ))
1284        .unwrap_err();
1285        assert!(err.to_string().contains("restore relay sequence overflow"));
1286        assert_eq!(
1287            block_on(store.get(&root(), &keys::sharding_marker())).unwrap(),
1288            None
1289        );
1290    }
1291
1292    #[test]
1293    fn single_namespace_epoch_rises_and_declares_recovery() {
1294        let snapshot = snapshot(
1295            &root(),
1296            vec![
1297                (keys::sharding_marker(), Value::new(b"single".to_vec())),
1298                (keys::grant_epoch(), codec::encode_u64(7)),
1299            ],
1300        );
1301        let store = MemoryKv::default();
1302        block_on(restore(
1303            &[snapshot],
1304            &store,
1305            RestoreOptions {
1306                epoch_at_least: None,
1307                recovered_at_ms: 44,
1308                allow_incomplete: false,
1309            },
1310        ))
1311        .unwrap();
1312        assert_eq!(
1313            block_on(store.get(&root(), &keys::grant_epoch())).unwrap(),
1314            Some(codec::encode_u64(7 + EPOCH_RESTORE_JUMP))
1315        );
1316        let lr = block_on(store.get(&root(), &keys::lease_recovery()))
1317            .unwrap()
1318            .unwrap();
1319        assert_eq!(codec::decode_lease_recovery(&lr).unwrap().resumed_at_ms, 44);
1320    }
1321
1322    #[test]
1323    fn absent_outbox_sequence_starts_next_allocation_above_target_watermark() {
1324        let root_snapshot = snapshot(
1325            &root(),
1326            vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1327        );
1328        let ref_snapshot = snapshot(
1329            &source(),
1330            vec![(
1331                keys::layout_version(),
1332                codec::encode_u32(keys::LAYOUT_VERSION),
1333            )],
1334        );
1335        let index_snapshot = snapshot(
1336            &target(),
1337            vec![(
1338                keys::relay_high_water(&source()).unwrap(),
1339                codec::encode_u64(5),
1340            )],
1341        );
1342        let store = MemoryKv::default();
1343        block_on(restore(
1344            &[ref_snapshot, index_snapshot, root_snapshot],
1345            &store,
1346            RestoreOptions {
1347                allow_incomplete: true,
1348                epoch_at_least: Some(EPOCH_RESTORE_JUMP + 99),
1349                ..RestoreOptions::default()
1350            },
1351        ))
1352        .unwrap();
1353        let os = block_on(store.get(&source(), &keys::outbox_sequence()))
1354            .unwrap()
1355            .unwrap();
1356        assert_eq!(codec::decode_u64(&os).unwrap(), 5);
1357        assert!(codec::decode_u64(&os).unwrap() + 1 > 5);
1358    }
1359
1360    #[test]
1361    fn missing_source_is_refused_or_reconstructed_above_target_watermark() {
1362        let archives = [
1363            snapshot(
1364                &root(),
1365                vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1366            ),
1367            snapshot(
1368                &coordinator(),
1369                vec![(keys::grant_epoch(), codec::encode_u64(1))],
1370            ),
1371            snapshot(
1372                &target(),
1373                vec![(
1374                    keys::relay_high_water(&source()).unwrap(),
1375                    codec::encode_u64(10),
1376                )],
1377            ),
1378        ];
1379        let store = RootFirst::default();
1380        let err = block_on(restore(&archives, &store, RestoreOptions::default())).unwrap_err();
1381        assert!(err.to_string().contains("missing relay sources"));
1382        assert!(err.to_string().contains("refs/heads/main"));
1383        let report = block_on(restore(
1384            &archives,
1385            &store,
1386            RestoreOptions {
1387                allow_incomplete: true,
1388                ..RestoreOptions::default()
1389            },
1390        ))
1391        .unwrap();
1392        assert_eq!(report.missing_sources, vec![source()]);
1393        assert_eq!(
1394            block_on(store.get(&source(), &keys::outbox_sequence())).unwrap(),
1395            Some(codec::encode_u64(10))
1396        );
1397        let relay = codec::encode_relay(&codec::RelayV1 {
1398            at_ms: 100,
1399            target: target(),
1400            puts: vec![(
1401                Key::new(b"m\0new".as_slice()),
1402                Value::new(b"delivered".as_slice()),
1403            )],
1404            deletes: Vec::new(),
1405        })
1406        .unwrap();
1407        block_on(
1408            store.apply(
1409                &source(),
1410                Batch::new()
1411                    .put(keys::outbox_sequence(), codec::encode_u64(11))
1412                    .put(keys::relay(11), relay),
1413            ),
1414        )
1415        .unwrap();
1416        let handler = RelayHandler {
1417            target: store,
1418            hook: NoHook,
1419            budget: RelayBudget::default(),
1420        };
1421        let ctx = TimerCtx {
1422            store: &handler.target,
1423            partition: &source(),
1424            now_ms: 200,
1425        };
1426        let timer = DueTimer {
1427            due_at_ms: 200,
1428            kind: kinds::RELAY,
1429            reference: bytes::Bytes::default(),
1430            value: Value::default(),
1431        };
1432        block_on(handler.fire(&ctx, &timer)).unwrap();
1433        assert_eq!(
1434            block_on(
1435                handler
1436                    .target
1437                    .get(&target(), &Key::new(b"m\0new".as_slice()))
1438            )
1439            .unwrap(),
1440            Some(Value::new(b"delivered".as_slice()))
1441        );
1442        assert_eq!(
1443            block_on(
1444                handler
1445                    .target
1446                    .get(&target(), &keys::relay_high_water(&source()).unwrap())
1447            )
1448            .unwrap(),
1449            Some(codec::encode_u64(11))
1450        );
1451    }
1452
1453    #[test]
1454    fn missing_coordinator_requires_both_flags_and_is_recovered_before_ref() {
1455        let archives = [
1456            snapshot(
1457                &root(),
1458                vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1459            ),
1460            snapshot(
1461                &source(),
1462                vec![(keys::outbox_sequence(), codec::encode_u64(1))],
1463            ),
1464        ];
1465        let store = RootFirst::default();
1466        for opts in [
1467            RestoreOptions::default(),
1468            RestoreOptions {
1469                allow_incomplete: true,
1470                ..RestoreOptions::default()
1471            },
1472        ] {
1473            let err = block_on(restore(&archives, &store, opts)).unwrap_err();
1474            assert!(err.to_string().contains("missing namespace coordinators"));
1475        }
1476        let low = block_on(restore(
1477            &archives,
1478            &store,
1479            RestoreOptions {
1480                allow_incomplete: true,
1481                epoch_at_least: Some(42),
1482                ..RestoreOptions::default()
1483            },
1484        ))
1485        .unwrap_err();
1486        assert!(low.to_string().contains("must be at least"));
1487        let floor = EPOCH_RESTORE_JUMP + 42;
1488        let report = block_on(restore(
1489            &archives,
1490            &store,
1491            RestoreOptions {
1492                allow_incomplete: true,
1493                epoch_at_least: Some(floor),
1494                ..RestoreOptions::default()
1495            },
1496        ))
1497        .unwrap();
1498        assert_eq!(report.missing_coordinators, vec![coordinator()]);
1499        assert_eq!(
1500            block_on(store.get(&coordinator(), &keys::grant_epoch())).unwrap(),
1501            Some(codec::encode_u64(floor))
1502        );
1503        assert!(
1504            block_on(store.get(&coordinator(), &keys::lease_recovery()))
1505                .unwrap()
1506                .is_some()
1507        );
1508        let writes = store.writes.lock().unwrap();
1509        assert!(
1510            writes.iter().position(|p| p == &coordinator()).unwrap()
1511                < writes.iter().position(|p| p == &source()).unwrap()
1512        );
1513    }
1514
1515    #[test]
1516    fn older_target_cannot_recover_already_delivered_rows() {
1517        let archives = [
1518            snapshot(
1519                &root(),
1520                vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1521            ),
1522            snapshot(
1523                &coordinator(),
1524                vec![(keys::grant_epoch(), codec::encode_u64(1))],
1525            ),
1526            snapshot(
1527                &source(),
1528                vec![(keys::outbox_sequence(), codec::encode_u64(5))],
1529            ),
1530            snapshot(
1531                &target(),
1532                vec![(
1533                    keys::relay_high_water(&source()).unwrap(),
1534                    codec::encode_u64(2),
1535                )],
1536            ),
1537        ];
1538        let store = MemoryKv::default();
1539        block_on(restore(&archives, &store, RestoreOptions::default())).unwrap();
1540        let os = block_on(store.get(&source(), &keys::outbox_sequence()))
1541            .unwrap()
1542            .unwrap();
1543        assert!(codec::decode_u64(&os).unwrap() > 2);
1544        assert_eq!(
1545            block_on(store.get(&target(), &Key::new(b"m\0already-delivered".as_slice()))).unwrap(),
1546            None
1547        );
1548    }
1549
1550    #[test]
1551    fn epoch_floor_and_overflow() {
1552        let root_archive = snapshot(
1553            &root(),
1554            vec![
1555                (keys::sharding_marker(), Value::new(b"single".to_vec())),
1556                (keys::grant_epoch(), codec::encode_u64(7)),
1557            ],
1558        );
1559        let store = MemoryKv::default();
1560        block_on(restore(
1561            &[root_archive],
1562            &store,
1563            RestoreOptions {
1564                epoch_at_least: Some(EPOCH_RESTORE_JUMP + 100),
1565                ..RestoreOptions::default()
1566            },
1567        ))
1568        .unwrap();
1569        assert_eq!(
1570            block_on(store.get(&root(), &keys::grant_epoch())).unwrap(),
1571            Some(codec::encode_u64(EPOCH_RESTORE_JUMP + 100))
1572        );
1573        let overflow = snapshot(
1574            &root(),
1575            vec![
1576                (keys::sharding_marker(), Value::new(b"single".to_vec())),
1577                (
1578                    keys::grant_epoch(),
1579                    codec::encode_u64(u64::MAX - EPOCH_RESTORE_JUMP + 1),
1580                ),
1581            ],
1582        );
1583        assert!(matches!(
1584            block_on(restore(
1585                &[overflow],
1586                &MemoryKv::default(),
1587                RestoreOptions::default()
1588            )),
1589            Err(StoreError::Invalid(_))
1590        ));
1591    }
1592    #[test]
1593    fn fenced_restore_preserves_mode_without_business_creation_and_rejects_missing_state() {
1594        let mode = codec::LeaseRecovery {
1595            authority_fence: Some(true),
1596            authority_ready: Some(true),
1597            activation_only: Some(true),
1598            resumed_at_ms: 0,
1599        };
1600        let archives = [
1601            snapshot(
1602                &root(),
1603                vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1604            ),
1605            snapshot(
1606                &coordinator(),
1607                vec![
1608                    (keys::authority_generation(), codec::encode_u64(7)),
1609                    (keys::lease_recovery(), codec::encode_lease_recovery(&mode)),
1610                ],
1611            ),
1612        ];
1613        let store = MemoryKv::default();
1614        block_on(restore(
1615            &archives,
1616            &store,
1617            RestoreOptions {
1618                recovered_at_ms: 100,
1619                ..RestoreOptions::default()
1620            },
1621        ))
1622        .unwrap();
1623        assert_eq!(
1624            block_on(store.get(&coordinator(), &keys::authority_generation())).unwrap(),
1625            Some(codec::encode_u64(7))
1626        );
1627        assert!(
1628            block_on(store.get(&coordinator(), &keys::namespace_record()))
1629                .unwrap()
1630                .is_none()
1631        );
1632        let raw = block_on(store.get(&coordinator(), &keys::lease_recovery()))
1633            .unwrap()
1634            .unwrap();
1635        let recovered = codec::decode_lease_recovery(&raw).unwrap();
1636        assert_eq!(recovered.authority_fence, Some(true));
1637        assert_eq!(recovered.authority_ready, Some(true));
1638        assert_eq!(recovered.recovery_time(), Some(100));
1639        let lease = codec::EpochLease {
1640            epoch: 0,
1641            config_version: 1,
1642            expires_at_ms: 1000,
1643            authority_generation: Some(7),
1644            authority_ready: Some(true),
1645        };
1646        let incomplete = [
1647            snapshot(
1648                &root(),
1649                vec![(keys::sharding_marker(), Value::new(b"d34".to_vec()))],
1650            ),
1651            snapshot(
1652                &source(),
1653                vec![(keys::epoch_lease(), codec::encode_epoch_lease(&lease))],
1654            ),
1655        ];
1656        let fresh = MemoryKv::default();
1657        assert!(
1658            block_on(restore(
1659                &incomplete,
1660                &fresh,
1661                RestoreOptions {
1662                    allow_incomplete: true,
1663                    epoch_at_least: Some(0),
1664                    ..RestoreOptions::default()
1665                }
1666            ))
1667            .is_err()
1668        );
1669        assert!(
1670            block_on(fresh.get(&root(), &keys::sharding_marker()))
1671                .unwrap()
1672                .is_none()
1673        );
1674    }
1675}