1use 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#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
34pub struct RestoreOptions {
35 pub epoch_at_least: Option<u64>,
37 pub recovered_at_ms: u64,
39 pub allow_incomplete: bool,
42}
43
44#[derive(Debug, Clone, Default, PartialEq, Eq)]
46pub struct RestoreReport {
47 pub partitions: usize,
49 pub records: u64,
51 pub missing_sources: Vec<Partition>,
53 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
300pub 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 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)] async 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 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 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
773pub 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}