1use std::fmt;
55use std::path::PathBuf;
56
57use mongreldb_types::hlc::HlcTimestamp;
58use mongreldb_types::ids::{NodeId, SchemaVersion, TableId, TabletId};
59use serde::{Deserialize, Serialize};
60use sha2::{Digest, Sha256};
61
62use crate::meta::MetaRejectionReason;
63use crate::node::ClusterError;
64use crate::split::{
65 ChildAllocation, ChildPlan, ChildProgress, ChildStateSink, SnapshotPin, SourceRetentionGuard,
66 TabletDataError, TabletKeyspace, TabletMetaPlane, TabletMutation,
67};
68use crate::tablet::{Key, ReplicaRole, TabletDescriptor, TabletError, TabletLayout, TabletState};
69
70#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
77pub enum MergeRejection {
78 #[error("tablets belong to different tables ({first_table} vs {second_table})")]
80 DifferentTables {
81 first_table: TableId,
83 second_table: TableId,
85 },
86 #[error("schema versions differ on table {table}: {first} vs {second}")]
88 SchemaMismatch {
89 table: TableId,
91 first: SchemaVersion,
93 second: SchemaVersion,
95 },
96 #[error("tablets {first} and {second} are not adjacent")]
98 NotAdjacent {
99 first: TabletId,
101 second: TabletId,
103 },
104 #[error("placement is incompatible: {first_nodes:?} vs {second_nodes:?}")]
108 IncompatiblePlacement {
109 first_nodes: Vec<NodeId>,
111 second_nodes: Vec<NodeId>,
113 },
114 #[error("schema job {job_id} is active on table {table}")]
117 ConflictingSchemaJob {
118 table: TableId,
120 job_id: u64,
122 },
123 #[error(
125 "combined size {combined_bytes} bytes exceeds the merge threshold {threshold_bytes} bytes"
126 )]
127 CombinedSizeExceedsThreshold {
128 combined_bytes: u64,
130 threshold_bytes: u64,
132 },
133 #[error("source tablet {tablet} is in state {state}, expected Active")]
135 InvalidSourceState {
136 tablet: TabletId,
138 state: TabletState,
140 },
141}
142
143#[derive(Debug, thiserror::Error)]
146pub enum MergeError {
147 #[error(transparent)]
149 Tablet(#[from] TabletError),
150 #[error(transparent)]
152 Fault(#[from] mongreldb_fault::Fault),
153 #[error(transparent)]
155 MetaPlane(#[from] MetaRejectionReason),
156 #[error(transparent)]
158 TabletData(#[from] TabletDataError),
159 #[error(transparent)]
161 Rejected(#[from] MergeRejection),
162 #[error("invalid merge plan: {0}")]
164 InvalidPlan(String),
165 #[error("applied key {0} lies outside its source partition")]
168 KeyOutsideSource(Key),
169 #[error("source tablet {tablet} is retained by {pins} old-generation pin(s)")]
171 SourceRetained {
172 tablet: TabletId,
174 pins: usize,
176 },
177}
178
179#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
186pub enum MergePhase {
187 Started,
189 MarkedMerging,
191 ReplacementCreated,
194 SnapshotsPinned,
196 ReplacementBuilt,
198 CaughtUp,
201 Published,
203 SourcesRetired,
205}
206
207impl MergePhase {
208 pub fn hook_name(self) -> Option<&'static str> {
212 Some(match self {
213 Self::Started => return None,
214 Self::MarkedMerging => "tablet.merge.phase.1",
215 Self::ReplacementCreated => "tablet.merge.phase.2",
216 Self::SnapshotsPinned => "tablet.merge.phase.3",
217 Self::ReplacementBuilt => "tablet.merge.phase.4",
218 Self::CaughtUp => "tablet.merge.phase.5",
219 Self::Published => "tablet.merge.phase.6",
220 Self::SourcesRetired => "tablet.merge.phase.7",
221 })
222 }
223}
224
225impl fmt::Display for MergePhase {
226 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
227 let name = match self {
228 Self::Started => "Started",
229 Self::MarkedMerging => "MarkedMerging",
230 Self::ReplacementCreated => "ReplacementCreated",
231 Self::SnapshotsPinned => "SnapshotsPinned",
232 Self::ReplacementBuilt => "ReplacementBuilt",
233 Self::CaughtUp => "CaughtUp",
234 Self::Published => "Published",
235 Self::SourcesRetired => "SourcesRetired",
236 };
237 f.write_str(name)
238 }
239}
240
241#[derive(Clone, Debug)]
248pub struct MergeInputs {
249 pub first: TabletDescriptor,
251 pub second: TabletDescriptor,
253 pub first_schema: SchemaVersion,
255 pub second_schema: SchemaVersion,
257 pub active_schema_job: Option<u64>,
260 pub first_size_bytes: u64,
262 pub second_size_bytes: u64,
264 pub max_merged_size_bytes: u64,
267}
268
269fn validate_merge_inputs(inputs: &MergeInputs) -> Result<[TabletDescriptor; 2], MergeRejection> {
272 let (first, second) = (&inputs.first, &inputs.second);
273 for source in [first, second] {
274 if source.state != TabletState::Active {
275 return Err(MergeRejection::InvalidSourceState {
276 tablet: source.tablet_id,
277 state: source.state,
278 });
279 }
280 }
281 if first.table_id != second.table_id {
282 return Err(MergeRejection::DifferentTables {
283 first_table: first.table_id,
284 second_table: second.table_id,
285 });
286 }
287 if inputs.first_schema != inputs.second_schema {
288 return Err(MergeRejection::SchemaMismatch {
289 table: first.table_id,
290 first: inputs.first_schema,
291 second: inputs.second_schema,
292 });
293 }
294 let ordered: [&TabletDescriptor; 2] = if first.partition.meets_start_of(&second.partition) {
295 [first, second]
296 } else if second.partition.meets_start_of(&first.partition) {
297 [second, first]
298 } else {
299 return Err(MergeRejection::NotAdjacent {
300 first: first.tablet_id,
301 second: second.tablet_id,
302 });
303 };
304 let nodes = |descriptor: &TabletDescriptor| -> Vec<NodeId> {
305 let mut nodes: Vec<NodeId> = descriptor
306 .replicas
307 .iter()
308 .map(|replica| replica.node_id)
309 .collect();
310 nodes.sort();
311 nodes
312 };
313 let (first_nodes, second_nodes) = (nodes(first), nodes(second));
314 if first_nodes != second_nodes {
315 return Err(MergeRejection::IncompatiblePlacement {
316 first_nodes,
317 second_nodes,
318 });
319 }
320 if let Some(job_id) = inputs.active_schema_job {
321 return Err(MergeRejection::ConflictingSchemaJob {
322 table: first.table_id,
323 job_id,
324 });
325 }
326 let combined = inputs
327 .first_size_bytes
328 .checked_add(inputs.second_size_bytes)
329 .ok_or(MergeRejection::CombinedSizeExceedsThreshold {
330 combined_bytes: u64::MAX,
331 threshold_bytes: inputs.max_merged_size_bytes,
332 })?;
333 if combined > inputs.max_merged_size_bytes {
334 return Err(MergeRejection::CombinedSizeExceedsThreshold {
335 combined_bytes: combined,
336 threshold_bytes: inputs.max_merged_size_bytes,
337 });
338 }
339 Ok([ordered[0].clone(), ordered[1].clone()])
340}
341
342#[derive(Clone, Debug)]
345pub struct MergePlan {
346 pub sources: [TabletDescriptor; 2],
349 pub replacement: ChildPlan,
352 pub merge_ts: HlcTimestamp,
354}
355
356impl MergePlan {
357 pub fn validate(&self) -> Result<(), MergeError> {
361 for source in &self.sources {
362 source.validate()?;
363 }
364 if !self.sources[0]
365 .partition
366 .meets_start_of(&self.sources[1].partition)
367 {
368 return Err(MergeError::InvalidPlan(
369 "merge sources are not ordered adjacent halves".to_owned(),
370 ));
371 }
372 let union = self.sources[0]
373 .partition
374 .union_adjacent(&self.sources[1].partition)
375 .ok_or_else(|| MergeError::InvalidPlan("source bounds are not adjacent".to_owned()))?;
376 if self.replacement.bounds != union {
377 return Err(MergeError::InvalidPlan(
378 "replacement bounds are not the union of the source bounds".to_owned(),
379 ));
380 }
381 if self.sources[0].tablet_id == self.sources[1].tablet_id
382 || self
383 .sources
384 .iter()
385 .any(|source| source.tablet_id == self.replacement.layout.tablet_id())
386 || self.sources[0].raft_group_id == self.sources[1].raft_group_id
387 || self
388 .sources
389 .iter()
390 .any(|source| source.raft_group_id == self.replacement.layout.raft_group_id())
391 {
392 return Err(MergeError::InvalidPlan(
393 "source and replacement tablet/raft-group ids must be distinct".to_owned(),
394 ));
395 }
396 self.replacement_descriptor().validate()?;
397 Ok(())
398 }
399
400 pub fn replacement_descriptor(&self) -> TabletDescriptor {
404 let generation = self
405 .sources
406 .iter()
407 .map(|source| source.generation)
408 .max()
409 .expect("two sources")
410 .checked_add(1)
411 .expect("descriptor generation overflows u64");
412 self.replacement.descriptor(
413 self.sources[0].table_id,
414 generation,
415 TabletState::Creating,
416 ReplicaRole::Learner,
417 )
418 }
419
420 pub fn publish_generation(&self) -> u64 {
422 self.sources
423 .iter()
424 .map(|source| source.generation)
425 .max()
426 .expect("two sources")
427 .checked_add(2)
428 .expect("descriptor generation overflows u64")
429 }
430}
431
432#[derive(Clone, Debug)]
434pub struct MergePlanner {
435 node_data: PathBuf,
436}
437
438impl MergePlanner {
439 pub fn new(node_data: impl Into<PathBuf>) -> Self {
441 Self {
442 node_data: node_data.into(),
443 }
444 }
445
446 pub fn plan(
449 &self,
450 inputs: MergeInputs,
451 merge_ts: HlcTimestamp,
452 allocation: ChildAllocation,
453 ) -> Result<MergePlan, MergeError> {
454 inputs.first.validate()?;
455 inputs.second.validate()?;
456 let sources = validate_merge_inputs(&inputs)?;
457 let bounds = sources[0]
458 .partition
459 .union_adjacent(&sources[1].partition)
460 .ok_or_else(|| MergeError::InvalidPlan("source bounds are not adjacent".to_owned()))?;
461 let plan = MergePlan {
462 sources,
463 replacement: ChildPlan {
464 bounds,
465 layout: TabletLayout::new(
466 self.node_data.clone(),
467 allocation.tablet_id,
468 allocation.raft_group_id,
469 ),
470 replicas: allocation.replicas,
471 },
472 merge_ts,
473 };
474 plan.validate()?;
475 Ok(plan)
476 }
477}
478
479#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
492pub struct MergePublishCommand {
493 pub sources: [TabletDescriptor; 2],
496 pub replacement: TabletDescriptor,
500 pub merge_ts: HlcTimestamp,
502}
503
504impl MergePublishCommand {
505 pub fn from_plan(plan: &MergePlan) -> Result<Self, MergeError> {
508 let publish_generation = plan.publish_generation();
509 let mut sources = plan.sources.clone();
510 for source in &mut sources {
511 let marked = source.published_transition(TabletState::Merging)?;
512 let mut retiring = marked.published_transition(TabletState::Retiring)?;
513 retiring.generation = publish_generation;
515 *source = retiring;
516 }
517 let mut replacement = plan
518 .replacement_descriptor()
519 .published_transition(TabletState::Active)?;
520 debug_assert_eq!(replacement.generation, publish_generation);
521 for replica in &mut replacement.replicas {
522 replica.role = ReplicaRole::Voter;
523 }
524 let command = Self {
525 sources,
526 replacement,
527 merge_ts: plan.merge_ts,
528 };
529 command.validate()?;
530 Ok(command)
531 }
532
533 pub fn publish_generation(&self) -> u64 {
535 self.replacement.generation
536 }
537
538 pub fn validate(&self) -> Result<(), MergeError> {
542 if self.replacement.state != TabletState::Active {
543 return Err(MergeError::InvalidPlan(format!(
544 "published replacement must be Active, is {}",
545 self.replacement.state
546 )));
547 }
548 let generation = self.replacement.generation;
549 for source in &self.sources {
550 if source.state != TabletState::Retiring {
551 return Err(MergeError::InvalidPlan(format!(
552 "published source {} must be Retiring, is {}",
553 source.tablet_id, source.state
554 )));
555 }
556 if source.generation != generation {
557 return Err(MergeError::InvalidPlan(
558 "publication assigns one generation to all descriptors".to_owned(),
559 ));
560 }
561 if source.table_id != self.replacement.table_id {
562 return Err(MergeError::InvalidPlan(
563 "sources and replacement name different tables".to_owned(),
564 ));
565 }
566 source.validate()?;
567 }
568 self.replacement.validate()?;
569 let union = self.sources[0]
570 .partition
571 .union_adjacent(&self.sources[1].partition)
572 .ok_or_else(|| MergeError::InvalidPlan("source bounds are not adjacent".to_owned()))?;
573 if self.replacement.partition != union {
574 return Err(MergeError::InvalidPlan(
575 "published replacement bounds are not the union of the source bounds".to_owned(),
576 ));
577 }
578 if self.sources[0].tablet_id == self.sources[1].tablet_id
579 || self
580 .sources
581 .iter()
582 .any(|source| source.tablet_id == self.replacement.tablet_id)
583 {
584 return Err(MergeError::InvalidPlan(
585 "publication tablet ids must be distinct".to_owned(),
586 ));
587 }
588 Ok(())
589 }
590}
591
592pub const MERGE_PROGRESS_FILENAME: &str = "merge.json";
598pub const MERGE_PROGRESS_FORMAT_VERSION: u32 = 1;
600pub const MIN_SUPPORTED_MERGE_PROGRESS_FORMAT_VERSION: u32 = 1;
602
603#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
606#[serde(deny_unknown_fields)]
607pub struct MergeProgress {
608 pub sources: [TabletDescriptor; 2],
611 pub replacement: ChildProgress,
613 pub merge_ts: HlcTimestamp,
615 pub phase: MergePhase,
617}
618
619impl MergeProgress {
620 pub fn from_plan(plan: &MergePlan, phase: MergePhase) -> Self {
622 Self {
623 sources: plan.sources.clone(),
624 replacement: plan.replacement.progress(),
625 merge_ts: plan.merge_ts,
626 phase,
627 }
628 }
629
630 pub fn plan(&self) -> MergePlan {
632 MergePlan {
633 sources: self.sources.clone(),
634 replacement: self.replacement.plan(),
635 merge_ts: self.merge_ts,
636 }
637 }
638}
639
640#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
643#[serde(deny_unknown_fields)]
644struct MergeProgressFile {
645 format_version: u32,
647 checksum: String,
649 progress: MergeProgress,
651}
652
653impl MergeProgressFile {
654 fn envelope(progress: &MergeProgress) -> Result<Self, MergeError> {
655 Ok(Self {
656 format_version: MERGE_PROGRESS_FORMAT_VERSION,
657 checksum: progress_checksum(progress).map_err(meta_io)?,
658 progress: progress.clone(),
659 })
660 }
661}
662
663fn progress_checksum(progress: &MergeProgress) -> Result<String, ClusterError> {
665 let bytes = serde_json::to_vec(progress).map_err(|error| ClusterError::CorruptMetadata {
666 file: MERGE_PROGRESS_FILENAME,
667 detail: format!("encode: {error}"),
668 })?;
669 Ok(hex_encode(&Sha256::digest(&bytes)))
670}
671
672fn meta_io(error: ClusterError) -> MergeError {
674 MergeError::Tablet(TabletError::Metadata(error))
675}
676
677fn hex_encode(bytes: &[u8]) -> String {
678 const HEX: &[u8; 16] = b"0123456789abcdef";
679 let mut out = String::with_capacity(bytes.len() * 2);
680 for byte in bytes {
681 out.push(HEX[(byte >> 4) as usize] as char);
682 out.push(HEX[(byte & 0x0f) as usize] as char);
683 }
684 out
685}
686
687pub trait MergeMetaPlane: TabletMetaPlane {
693 fn publish_merge(&mut self, command: &MergePublishCommand) -> Result<(), MetaRejectionReason> {
699 command
700 .validate()
701 .map_err(|error| MetaRejectionReason::Invalid {
702 reason: error.to_string(),
703 })?;
704 self.set_tablet(&command.replacement)?;
705 for source in &command.sources {
706 self.set_tablet(source)?;
707 }
708 Ok(())
709 }
710}
711
712impl MergeMetaPlane for crate::split::InMemoryMetaPlane {}
713
714pub struct MergeExecutor<M, K, S> {
727 progress: MergeProgress,
728 source_layouts: [TabletLayout; 2],
729 meta: M,
730 keyspaces: [K; 2],
731 sink: S,
732 snapshot_pins: [Option<Box<dyn SnapshotPin>>; 2],
733 retention: Option<[SourceRetentionGuard; 2]>,
734}
735
736impl<M: MergeMetaPlane, K: TabletKeyspace, S: ChildStateSink> MergeExecutor<M, K, S> {
737 pub fn begin(
741 plan: MergePlan,
742 source_layouts: [TabletLayout; 2],
743 meta: M,
744 keyspaces: [K; 2],
745 sink: S,
746 ) -> Result<Self, MergeError> {
747 plan.validate()?;
748 for (source, layout) in plan.sources.iter().zip(source_layouts.iter()) {
749 if source.state != TabletState::Active {
750 return Err(MergeRejection::InvalidSourceState {
751 tablet: source.tablet_id,
752 state: source.state,
753 }
754 .into());
755 }
756 if layout.tablet_id() != source.tablet_id
757 || layout.raft_group_id() != source.raft_group_id
758 {
759 return Err(TabletError::TabletMismatch {
760 path: layout.tablet_dir(),
761 expected: layout.tablet_id(),
762 found: source.tablet_id,
763 expected_group: layout.raft_group_id(),
764 found_group: source.raft_group_id,
765 }
766 .into());
767 }
768 layout.validate()?;
769 }
770 let executor = Self {
771 progress: MergeProgress::from_plan(&plan, MergePhase::Started),
772 source_layouts,
773 meta,
774 keyspaces,
775 sink,
776 snapshot_pins: [None, None],
777 retention: None,
778 };
779 executor.persist_progress()?;
780 Ok(executor)
781 }
782
783 pub fn resume(
788 source_layouts: [TabletLayout; 2],
789 meta: M,
790 keyspaces: [K; 2],
791 sink: S,
792 ) -> Result<Option<Self>, MergeError> {
793 let Some(progress) = load_progress(&source_layouts[0])? else {
794 return Ok(None);
795 };
796 for (source, layout) in progress.sources.iter().zip(source_layouts.iter()) {
797 if source.tablet_id != layout.tablet_id()
798 || source.raft_group_id != layout.raft_group_id()
799 {
800 return Err(TabletError::TabletMismatch {
801 path: layout.tablet_dir(),
802 expected: layout.tablet_id(),
803 found: source.tablet_id,
804 expected_group: layout.raft_group_id(),
805 found_group: source.raft_group_id,
806 }
807 .into());
808 }
809 }
810 progress.plan().validate()?;
811 Ok(Some(Self {
812 progress,
813 source_layouts,
814 meta,
815 keyspaces,
816 sink,
817 snapshot_pins: [None, None],
818 retention: None,
819 }))
820 }
821
822 pub fn phase(&self) -> MergePhase {
824 self.progress.phase
825 }
826
827 pub fn progress(&self) -> &MergeProgress {
829 &self.progress
830 }
831
832 pub fn plan(&self) -> MergePlan {
834 self.progress.plan()
835 }
836
837 pub fn retention(&self) -> Option<&[SourceRetentionGuard; 2]> {
839 self.retention.as_ref()
840 }
841
842 pub fn meta(&self) -> &M {
844 &self.meta
845 }
846
847 pub fn keyspaces(&self) -> &[K; 2] {
849 &self.keyspaces
850 }
851
852 pub fn sink(&self) -> &S {
854 &self.sink
855 }
856
857 pub fn step(&mut self) -> Result<MergePhase, MergeError> {
861 use MergePhase::{
862 CaughtUp, MarkedMerging, Published, ReplacementBuilt, ReplacementCreated,
863 SnapshotsPinned, SourcesRetired, Started,
864 };
865 let next = match self.progress.phase {
866 Started => {
867 self.mark_sources_merging()?;
868 MarkedMerging
869 }
870 MarkedMerging => {
871 self.create_replacement()?;
872 ReplacementCreated
873 }
874 ReplacementCreated => {
875 self.pin_source_snapshots()?;
876 SnapshotsPinned
877 }
878 SnapshotsPinned => {
879 self.build_replacement()?;
880 ReplacementBuilt
881 }
882 ReplacementBuilt => {
883 self.catch_up_replacement()?;
884 CaughtUp
885 }
886 CaughtUp => {
887 self.publish_replacement()?;
888 Published
889 }
890 Published => {
891 self.remove_sources()?;
892 SourcesRetired
893 }
894 SourcesRetired => SourcesRetired,
895 };
896 self.progress.phase = next;
899 if next != SourcesRetired {
900 self.persist_progress()?;
901 }
902 if let Some(hook) = next.hook_name() {
903 mongreldb_fault::inject(hook)?;
904 }
905 Ok(next)
906 }
907
908 pub fn run(&mut self) -> Result<(), MergeError> {
913 while self.progress.phase != MergePhase::SourcesRetired {
914 self.step()?;
915 }
916 Ok(())
917 }
918
919 pub fn run_until(&mut self, phase: MergePhase) -> Result<(), MergeError> {
921 while self.progress.phase != phase {
922 self.step()?;
923 }
924 Ok(())
925 }
926
927 fn mark_sources_merging(&mut self) -> Result<(), MergeError> {
930 for (source, layout) in self.progress.sources.iter().zip(self.source_layouts.iter()) {
931 let marked = source.published_transition(TabletState::Merging)?;
932 self.meta.set_tablet(&marked)?;
933 layout.store_metadata(&marked)?;
934 }
935 Ok(())
937 }
938
939 fn create_replacement(&mut self) -> Result<(), MergeError> {
942 let plan = self.plan();
943 let descriptor = plan.replacement_descriptor();
944 plan.replacement.layout.create(&descriptor)?;
945 self.meta.set_tablet(&descriptor)?;
946 Ok(())
947 }
948
949 fn pin_source_snapshots(&mut self) -> Result<(), MergeError> {
952 self.ensure_snapshot_pins()
953 }
954
955 fn ensure_snapshot_pins(&mut self) -> Result<(), MergeError> {
957 for index in 0..2 {
958 if self.snapshot_pins[index].is_none() {
959 self.snapshot_pins[index] =
960 Some(self.keyspaces[index].pin_snapshot(self.progress.merge_ts)?);
961 }
962 }
963 Ok(())
964 }
965
966 fn build_replacement(&mut self) -> Result<(), MergeError> {
971 self.ensure_snapshot_pins()?;
972 self.sink.begin_build()?;
973 for index in 0..2 {
974 let snapshot = self.keyspaces[index].snapshot_at(self.progress.merge_ts)?;
975 for (key, value) in snapshot {
976 if !self.progress.sources[index].partition.contains(&key) {
977 return Err(MergeError::KeyOutsideSource(key));
978 }
979 self.sink.stage(&key, &value)?;
980 }
981 }
982 self.sink.install_staged()?;
983 Ok(())
984 }
985
986 fn catch_up_replacement(&mut self) -> Result<(), MergeError> {
989 self.ensure_snapshot_pins()?;
990 for index in 0..2 {
991 let deltas = self.keyspaces[index].deltas_after(self.progress.merge_ts)?;
992 for mutation in deltas {
993 let key = match &mutation {
994 TabletMutation::Upsert(key, _) | TabletMutation::Delete(key) => key,
995 };
996 if !self.progress.sources[index].partition.contains(key) {
997 return Err(MergeError::KeyOutsideSource(key.clone()));
998 }
999 self.sink.apply_delta(&mutation)?;
1000 }
1001 }
1002 Ok(())
1003 }
1004
1005 fn publish_replacement(&mut self) -> Result<(), MergeError> {
1009 let plan = self.plan();
1010 let command = MergePublishCommand::from_plan(&plan)?;
1011 mongreldb_fault::inject("tablet.merge.before")?;
1012 self.meta.publish_merge(&command)?;
1013 mongreldb_fault::inject("tablet.merge.after")?;
1014 plan.replacement
1016 .layout
1017 .store_metadata(&command.replacement)?;
1018 for (descriptor, layout) in command.sources.iter().zip(self.source_layouts.iter()) {
1019 layout.store_metadata(descriptor)?;
1020 }
1021 self.retention = Some(
1022 command
1023 .sources
1024 .clone()
1025 .map(|source| SourceRetentionGuard::new(source.tablet_id, source.generation)),
1026 );
1027 self.snapshot_pins = [None, None];
1029 Ok(())
1030 }
1031
1032 fn remove_sources(&mut self) -> Result<(), MergeError> {
1036 if let Some(guards) = &self.retention {
1037 for guard in guards {
1038 if !guard.ready_for_removal() {
1039 return Err(MergeError::SourceRetained {
1040 tablet: guard.source(),
1041 pins: guard.old_generation_pins(),
1042 });
1043 }
1044 }
1045 }
1046 for index in 0..2 {
1047 let source_id = self.progress.sources[index].tablet_id;
1048 if let Some(current) = self.meta.tablet(source_id) {
1049 let retired = if current.state == TabletState::Retired {
1050 current
1051 } else {
1052 let retired = current.published_transition(TabletState::Retired)?;
1053 self.meta.set_tablet(&retired)?;
1054 retired
1055 };
1056 self.meta.remove_tablet(source_id, retired.generation)?;
1057 }
1058 }
1059 self.source_layouts[1].teardown()?;
1060 self.source_layouts[0].teardown()?;
1061 Ok(())
1062 }
1063
1064 fn persist_progress(&self) -> Result<(), MergeError> {
1067 let file = MergeProgressFile::envelope(&self.progress)?;
1068 let bytes = crate::node::encode_json(MERGE_PROGRESS_FILENAME, &file).map_err(meta_io)?;
1069 crate::node::write_meta_atomic(
1070 &self.source_layouts[0].tablet_dir(),
1071 MERGE_PROGRESS_FILENAME,
1072 &bytes,
1073 )
1074 .map_err(ClusterError::Io)
1075 .map_err(meta_io)?;
1076 Ok(())
1077 }
1078}
1079
1080fn load_progress(lower_layout: &TabletLayout) -> Result<Option<MergeProgress>, MergeError> {
1083 let path = lower_layout.tablet_dir().join(MERGE_PROGRESS_FILENAME);
1084 let Some(bytes) = crate::node::read_meta_file(&path).map_err(meta_io)? else {
1085 return Ok(None);
1086 };
1087 let file: MergeProgressFile =
1088 crate::node::decode_json(MERGE_PROGRESS_FILENAME, &bytes).map_err(meta_io)?;
1089 if file.format_version < MIN_SUPPORTED_MERGE_PROGRESS_FORMAT_VERSION
1090 || file.format_version > MERGE_PROGRESS_FORMAT_VERSION
1091 {
1092 return Err(meta_io(ClusterError::UnsupportedFormatVersion {
1093 file: MERGE_PROGRESS_FILENAME,
1094 found: file.format_version,
1095 min: MIN_SUPPORTED_MERGE_PROGRESS_FORMAT_VERSION,
1096 max: MERGE_PROGRESS_FORMAT_VERSION,
1097 }));
1098 }
1099 if file.checksum != progress_checksum(&file.progress).map_err(meta_io)? {
1100 return Err(meta_io(ClusterError::CorruptMetadata {
1101 file: MERGE_PROGRESS_FILENAME,
1102 detail: "checksum mismatch".to_owned(),
1103 }));
1104 }
1105 let progress = file.progress;
1106 if progress.sources[0].tablet_id != lower_layout.tablet_id()
1107 || progress.sources[0].raft_group_id != lower_layout.raft_group_id()
1108 {
1109 return Err(TabletError::TabletMismatch {
1110 path: lower_layout.tablet_dir(),
1111 expected: lower_layout.tablet_id(),
1112 found: progress.sources[0].tablet_id,
1113 expected_group: lower_layout.raft_group_id(),
1114 found_group: progress.sources[0].raft_group_id,
1115 }
1116 .into());
1117 }
1118 Ok(Some(progress))
1119}
1120
1121pub fn merge_progress(lower_layout: &TabletLayout) -> Result<Option<MergeProgress>, MergeError> {
1125 load_progress(lower_layout)
1126}
1127
1128#[cfg(test)]
1133mod tests {
1134 use mongreldb_types::ids::RaftGroupId;
1135
1136 use super::*;
1137 use crate::split::{
1138 retry_guidance, ChildAllocation, InMemoryMetaPlane, MapChildSink, MapKeyspace,
1139 RetryGuidance, EXECUTOR_TEST_LOCK,
1140 };
1141 use crate::tablet::{
1142 check_generation, find_tablet_for_key, tablets_overlapping, Bound, KeyValue,
1143 PartitionBounds, ReplicaDescriptor, RoutingError, RowKeyEncoder,
1144 };
1145
1146 fn node(byte: u8) -> NodeId {
1147 NodeId::from_bytes([byte; 16])
1148 }
1149
1150 fn tablet_id(byte: u8) -> TabletId {
1151 TabletId::from_bytes([byte; 16])
1152 }
1153
1154 fn group_id(byte: u8) -> RaftGroupId {
1155 RaftGroupId::from_bytes([byte; 16])
1156 }
1157
1158 fn text_key(text: &str) -> Key {
1159 RowKeyEncoder::encode_key(&[KeyValue::Text(text.to_owned())])
1160 }
1161
1162 fn ts(micros: u64) -> HlcTimestamp {
1163 HlcTimestamp {
1164 physical_micros: micros,
1165 logical: 0,
1166 node_tiebreaker: 0,
1167 }
1168 }
1169
1170 fn voters() -> Vec<ReplicaDescriptor> {
1171 vec![
1172 ReplicaDescriptor {
1173 node_id: node(1),
1174 role: ReplicaRole::Voter,
1175 raft_node_id: 11,
1176 },
1177 ReplicaDescriptor {
1178 node_id: node(2),
1179 role: ReplicaRole::Voter,
1180 raft_node_id: 12,
1181 },
1182 ]
1183 }
1184
1185 fn source(
1186 tablet: u8,
1187 group: u8,
1188 low: Bound<Key>,
1189 high: Bound<Key>,
1190 generation: u64,
1191 ) -> TabletDescriptor {
1192 TabletDescriptor {
1193 tablet_id: tablet_id(tablet),
1194 table_id: TableId::new(3),
1195 database_id: mongreldb_types::ids::DatabaseId::ZERO,
1196 raft_group_id: group_id(group),
1197 partition: PartitionBounds::new(low, high).unwrap(),
1198 replicas: voters(),
1199 leader_hint: Some(node(1)),
1200 generation,
1201 state: TabletState::Active,
1202 }
1203 }
1204
1205 fn source_pair() -> (TabletDescriptor, TabletDescriptor) {
1208 (
1209 source(
1210 11,
1211 11,
1212 Bound::Included(text_key("a")),
1213 Bound::Excluded(text_key("m")),
1214 4,
1215 ),
1216 source(
1217 12,
1218 12,
1219 Bound::Included(text_key("m")),
1220 Bound::Excluded(text_key("z")),
1221 6,
1222 ),
1223 )
1224 }
1225
1226 fn inputs(first: TabletDescriptor, second: TabletDescriptor) -> MergeInputs {
1227 MergeInputs {
1228 first,
1229 second,
1230 first_schema: SchemaVersion::new(9),
1231 second_schema: SchemaVersion::new(9),
1232 active_schema_job: None,
1233 first_size_bytes: 1_000,
1234 second_size_bytes: 2_000,
1235 max_merged_size_bytes: 1 << 20,
1236 }
1237 }
1238
1239 fn replacement_allocation() -> ChildAllocation {
1240 ChildAllocation {
1241 tablet_id: tablet_id(13),
1242 raft_group_id: group_id(13),
1243 replicas: vec![
1244 ReplicaDescriptor {
1245 node_id: node(3),
1246 role: ReplicaRole::Voter,
1247 raft_node_id: 41,
1248 },
1249 ReplicaDescriptor {
1250 node_id: node(4),
1251 role: ReplicaRole::Voter,
1252 raft_node_id: 42,
1253 },
1254 ],
1255 }
1256 }
1257
1258 #[test]
1261 fn merge_validation_rejects_each_violated_requirement() {
1262 let dir = tempfile::tempdir().unwrap();
1263 let planner = MergePlanner::new(dir.path());
1264 let (left, right) = source_pair();
1265 let plan = |inputs: MergeInputs| planner.plan(inputs, ts(150), replacement_allocation());
1266
1267 let merged = plan(inputs(left.clone(), right.clone())).unwrap();
1269 assert_eq!(merged.sources[0].tablet_id, left.tablet_id);
1270 assert_eq!(merged.sources[1].tablet_id, right.tablet_id);
1271 assert_eq!(
1272 merged.replacement.bounds,
1273 PartitionBounds::new(
1274 Bound::Included(text_key("a")),
1275 Bound::Excluded(text_key("z"))
1276 )
1277 .unwrap()
1278 );
1279 let swapped = plan(inputs(right.clone(), left.clone())).unwrap();
1281 assert_eq!(swapped.sources[0].tablet_id, left.tablet_id);
1282
1283 let mut foreign = right.clone();
1285 foreign.table_id = TableId::new(4);
1286 assert!(matches!(
1287 plan(inputs(left.clone(), foreign)),
1288 Err(MergeError::Rejected(MergeRejection::DifferentTables {
1289 first_table,
1290 second_table,
1291 })) if first_table == TableId::new(3) && second_table == TableId::new(4)
1292 ));
1293
1294 let mut mismatched = inputs(left.clone(), right.clone());
1296 mismatched.second_schema = SchemaVersion::new(10);
1297 assert!(matches!(
1298 plan(mismatched),
1299 Err(MergeError::Rejected(MergeRejection::SchemaMismatch {
1300 table,
1301 first,
1302 second,
1303 })) if table == TableId::new(3)
1304 && first == SchemaVersion::new(9)
1305 && second == SchemaVersion::new(10)
1306 ));
1307
1308 let gapped = source(
1310 14,
1311 14,
1312 Bound::Included(text_key("n")),
1313 Bound::Excluded(text_key("z")),
1314 6,
1315 );
1316 assert!(matches!(
1317 plan(inputs(left.clone(), gapped)),
1318 Err(MergeError::Rejected(MergeRejection::NotAdjacent { .. }))
1319 ));
1320 let overlapping = source(
1321 14,
1322 14,
1323 Bound::Included(text_key("l")),
1324 Bound::Excluded(text_key("z")),
1325 6,
1326 );
1327 assert!(matches!(
1328 plan(inputs(left.clone(), overlapping)),
1329 Err(MergeError::Rejected(MergeRejection::NotAdjacent { .. }))
1330 ));
1331
1332 let mut elsewhere = right.clone();
1334 elsewhere.replicas[0].node_id = node(9);
1335 elsewhere.leader_hint = Some(node(9));
1336 assert!(matches!(
1337 plan(inputs(left.clone(), elsewhere)),
1338 Err(MergeError::Rejected(
1339 MergeRejection::IncompatiblePlacement { .. }
1340 ))
1341 ));
1342
1343 let mut with_job = inputs(left.clone(), right.clone());
1345 with_job.active_schema_job = Some(77);
1346 assert!(matches!(
1347 plan(with_job),
1348 Err(MergeError::Rejected(MergeRejection::ConflictingSchemaJob {
1349 table,
1350 job_id: 77,
1351 })) if table == TableId::new(3)
1352 ));
1353
1354 let mut too_big = inputs(left.clone(), right.clone());
1356 too_big.max_merged_size_bytes = 2_999;
1357 assert!(matches!(
1358 plan(too_big),
1359 Err(MergeError::Rejected(
1360 MergeRejection::CombinedSizeExceedsThreshold {
1361 combined_bytes: 3_000,
1362 threshold_bytes: 2_999,
1363 }
1364 ))
1365 ));
1366 let mut overflowing = inputs(left.clone(), right.clone());
1367 overflowing.first_size_bytes = u64::MAX;
1368 assert!(matches!(
1369 plan(overflowing),
1370 Err(MergeError::Rejected(
1371 MergeRejection::CombinedSizeExceedsThreshold { .. }
1372 ))
1373 ));
1374
1375 let mut splitting = left.clone();
1377 splitting.state = TabletState::Splitting;
1378 assert!(matches!(
1379 plan(inputs(splitting, right.clone())),
1380 Err(MergeError::Rejected(MergeRejection::InvalidSourceState {
1381 state: TabletState::Splitting,
1382 ..
1383 }))
1384 ));
1385
1386 let mut colliding = replacement_allocation();
1388 colliding.tablet_id = left.tablet_id;
1389 assert!(matches!(
1390 planner.plan(inputs(left.clone(), right.clone()), ts(150), colliding),
1391 Err(MergeError::InvalidPlan(_))
1392 ));
1393 }
1394
1395 #[test]
1396 fn merge_publish_command_flips_three_descriptors_at_one_generation() {
1397 let dir = tempfile::tempdir().unwrap();
1398 let planner = MergePlanner::new(dir.path());
1399 let (left, right) = source_pair();
1400 let plan = planner
1401 .plan(inputs(left, right), ts(150), replacement_allocation())
1402 .unwrap();
1403 assert_eq!(plan.replacement_descriptor().generation, 7);
1405 let command = MergePublishCommand::from_plan(&plan).unwrap();
1406 assert_eq!(command.publish_generation(), 8);
1407 assert_eq!(command.replacement.state, TabletState::Active);
1408 assert!(command
1409 .replacement
1410 .replicas
1411 .iter()
1412 .all(|replica| replica.role == ReplicaRole::Voter));
1413 for (source, original) in command.sources.iter().zip(plan.sources.iter()) {
1414 assert_eq!(source.state, TabletState::Retiring);
1415 assert_eq!(source.generation, 8, "source at g={}", original.generation);
1416 }
1417 let bytes = serde_json::to_vec(&command).unwrap();
1419 let back: MergePublishCommand = serde_json::from_slice(&bytes).unwrap();
1420 assert_eq!(back, command);
1421
1422 let mut wrong_state = command.clone();
1424 wrong_state.replacement.state = TabletState::Creating;
1425 assert!(wrong_state.validate().is_err());
1426 let mut skewed = command.clone();
1427 skewed.sources[0].generation = 9;
1428 assert!(skewed.validate().is_err());
1429 let mut wrong_bounds = command.clone();
1430 wrong_bounds.replacement.partition =
1431 PartitionBounds::new(Bound::Unbounded, Bound::Unbounded).unwrap();
1432 assert!(wrong_bounds.validate().is_err());
1433 }
1434
1435 struct MergeFixture {
1438 _dir: tempfile::TempDir,
1439 sources: [TabletDescriptor; 2],
1440 source_layouts: [TabletLayout; 2],
1441 meta: InMemoryMetaPlane,
1442 keyspaces: [MapKeyspace; 2],
1443 sink: MapChildSink,
1444 plan: MergePlan,
1445 }
1446
1447 const LOWER_KEYS: [&str; 6] = ["b", "d", "f", "h", "j", "l"];
1448 const UPPER_KEYS: [&str; 7] = ["m", "o", "q", "s", "u", "w", "y"];
1449
1450 fn merge_fixture() -> MergeFixture {
1451 let dir = tempfile::tempdir().unwrap();
1452 let (left, right) = source_pair();
1453 let mut meta = InMemoryMetaPlane::new();
1454 meta.set_tablet(&left).unwrap();
1455 meta.set_tablet(&right).unwrap();
1456
1457 let left_layout = TabletLayout::new(dir.path(), left.tablet_id, left.raft_group_id);
1458 left_layout.create(&left).unwrap();
1459 let right_layout = TabletLayout::new(dir.path(), right.tablet_id, right.raft_group_id);
1460 right_layout.create(&right).unwrap();
1461
1462 let left_keyspace = MapKeyspace::new();
1463 for name in LOWER_KEYS {
1464 left_keyspace.insert(
1465 text_key(name),
1466 ts(100),
1467 format!("v-{name}@100").into_bytes(),
1468 );
1469 }
1470 left_keyspace.insert(text_key("d"), ts(200), b"v-d@200".to_vec());
1471 let right_keyspace = MapKeyspace::new();
1472 for name in UPPER_KEYS {
1473 right_keyspace.insert(
1474 text_key(name),
1475 ts(100),
1476 format!("v-{name}@100").into_bytes(),
1477 );
1478 }
1479 right_keyspace.insert(text_key("x"), ts(200), b"v-x@200".to_vec());
1480
1481 let planner = MergePlanner::new(dir.path());
1482 let plan = planner
1483 .plan(
1484 inputs(left.clone(), right.clone()),
1485 ts(150),
1486 replacement_allocation(),
1487 )
1488 .unwrap();
1489 MergeFixture {
1490 _dir: dir,
1491 sources: [left, right],
1492 source_layouts: [left_layout, right_layout],
1493 meta,
1494 keyspaces: [left_keyspace, right_keyspace],
1495 sink: MapChildSink::new(),
1496 plan,
1497 }
1498 }
1499
1500 type TestExecutor = MergeExecutor<InMemoryMetaPlane, MapKeyspace, MapChildSink>;
1501
1502 fn begin_executor(fixture: &MergeFixture) -> TestExecutor {
1503 MergeExecutor::begin(
1504 fixture.plan.clone(),
1505 fixture.source_layouts.clone(),
1506 fixture.meta.clone(),
1507 fixture.keyspaces.clone(),
1508 fixture.sink.clone(),
1509 )
1510 .unwrap()
1511 }
1512
1513 fn assert_merge_completed(fixture: &MergeFixture) {
1514 for source in &fixture.sources {
1517 assert!(fixture.meta.tablet(source.tablet_id).is_none());
1518 }
1519 let replacement = fixture.meta.tablet(tablet_id(13)).unwrap();
1520 assert_eq!(replacement.state, TabletState::Active);
1521 assert_eq!(replacement.generation, 8);
1522 assert!(replacement
1523 .replicas
1524 .iter()
1525 .all(|replica| replica.role == ReplicaRole::Voter));
1526 assert_eq!(replacement.partition, fixture.plan.replacement.bounds);
1527 assert_eq!(
1528 fixture.plan.replacement.layout.load_metadata().unwrap(),
1529 replacement
1530 );
1531 for layout in &fixture.source_layouts {
1533 assert!(!layout.tablet_dir().exists());
1534 assert!(!layout.group_dir().exists());
1535 }
1536 assert!(!fixture.source_layouts[0]
1537 .tablet_dir()
1538 .join(MERGE_PROGRESS_FILENAME)
1539 .exists());
1540 let rows = fixture.sink.rows();
1542 let mut expected = fixture.keyspaces[0].rows_at(ts(u64::MAX));
1543 expected.extend(fixture.keyspaces[1].rows_at(ts(u64::MAX)));
1544 assert_eq!(rows, expected);
1545 }
1546
1547 #[test]
1548 fn full_merge_replaces_adjacent_tablets_with_zero_loss() {
1549 let _lock = EXECUTOR_TEST_LOCK.lock().unwrap();
1550 let fixture = merge_fixture();
1551 let table = TableId::new(3);
1552 let mut executor = begin_executor(&fixture);
1553 assert_eq!(executor.phase(), MergePhase::Started);
1554
1555 assert_eq!(executor.step().unwrap(), MergePhase::MarkedMerging);
1557 let marked_left = fixture.meta.tablet(tablet_id(11)).unwrap();
1558 let marked_right = fixture.meta.tablet(tablet_id(12)).unwrap();
1559 assert_eq!(marked_left.state, TabletState::Merging);
1560 assert_eq!(marked_left.generation, 5);
1561 assert_eq!(marked_right.state, TabletState::Merging);
1562 assert_eq!(marked_right.generation, 7);
1563 let error = check_generation(&marked_left, 4).unwrap_err();
1566 assert!(matches!(error, RoutingError::StaleMetadata { .. }));
1567 let tablets = fixture.meta.descriptors();
1568 assert_eq!(
1569 find_tablet_for_key(&tablets, table, &text_key("b"))
1570 .unwrap()
1571 .tablet_id,
1572 tablet_id(11)
1573 );
1574 assert_eq!(
1575 find_tablet_for_key(&tablets, table, &text_key("y"))
1576 .unwrap()
1577 .tablet_id,
1578 tablet_id(12)
1579 );
1580
1581 assert_eq!(executor.step().unwrap(), MergePhase::ReplacementCreated);
1583 let creating = fixture.meta.tablet(tablet_id(13)).unwrap();
1584 assert_eq!(creating.state, TabletState::Creating);
1585 assert_eq!(creating.generation, 7);
1586 assert!(creating
1587 .replicas
1588 .iter()
1589 .all(|replica| replica.role == ReplicaRole::Learner));
1590 let tablets = fixture.meta.descriptors();
1591 assert_eq!(
1592 find_tablet_for_key(&tablets, table, &text_key("b"))
1593 .unwrap()
1594 .tablet_id,
1595 tablet_id(11),
1596 "hidden replacement exposed before catch-up"
1597 );
1598
1599 assert_eq!(executor.step().unwrap(), MergePhase::SnapshotsPinned);
1601 assert_eq!(fixture.keyspaces[0].pin_count(), 1);
1602 assert_eq!(fixture.keyspaces[1].pin_count(), 1);
1603 assert_eq!(executor.step().unwrap(), MergePhase::ReplacementBuilt);
1604 assert_eq!(
1605 fixture.sink.rows().len(),
1606 LOWER_KEYS.len() + UPPER_KEYS.len()
1607 );
1608 assert_eq!(
1609 fixture.sink.rows().get(&text_key("d")),
1610 Some(&b"v-d@100".to_vec()),
1611 "post-merge write leaked into the pinned snapshot"
1612 );
1613 assert_eq!(executor.step().unwrap(), MergePhase::CaughtUp);
1614 assert_eq!(
1615 fixture.sink.rows().get(&text_key("d")),
1616 Some(&b"v-d@200".to_vec())
1617 );
1618 assert_eq!(
1619 fixture.sink.rows().get(&text_key("x")),
1620 Some(&b"v-x@200".to_vec())
1621 );
1622
1623 assert_eq!(executor.step().unwrap(), MergePhase::Published);
1625 assert_eq!(fixture.keyspaces[0].pin_count(), 0);
1626 assert_eq!(fixture.keyspaces[1].pin_count(), 0);
1627 for id in [tablet_id(11), tablet_id(12)] {
1628 let retiring = fixture.meta.tablet(id).unwrap();
1629 assert_eq!(retiring.state, TabletState::Retiring);
1630 assert_eq!(retiring.generation, 8);
1631 }
1632 let tablets = fixture.meta.descriptors();
1633 for name in ["b", "m", "y"] {
1634 assert_eq!(
1635 find_tablet_for_key(&tablets, table, &text_key(name))
1636 .unwrap()
1637 .tablet_id,
1638 tablet_id(13),
1639 "key {name} did not reroute to the replacement"
1640 );
1641 }
1642 let overlapping = tablets_overlapping(&tablets, table, &PartitionBounds::unbounded());
1643 assert_eq!(
1644 overlapping
1645 .iter()
1646 .map(|tablet| tablet.tablet_id)
1647 .collect::<Vec<_>>(),
1648 vec![tablet_id(13)]
1649 );
1650 let retiring = fixture.meta.tablet(tablet_id(11)).unwrap();
1652 let error = check_generation(&retiring, 4).unwrap_err();
1653 assert!(matches!(error, RoutingError::TabletMoved { .. }));
1654 assert!(matches!(
1655 retry_guidance(&error),
1656 RetryGuidance::RefreshAndReroute { .. }
1657 ));
1658 assert!(check_generation(&fixture.meta.tablet(tablet_id(13)).unwrap(), 8).is_ok());
1659
1660 let pin_left = executor.retention().unwrap()[0].pin(4);
1662 assert!(matches!(
1663 executor.step(),
1664 Err(MergeError::SourceRetained { pins: 1, .. })
1665 ));
1666 assert_eq!(executor.phase(), MergePhase::Published);
1667 assert!(executor.retention().unwrap()[0].unpin(pin_left));
1668 let pin_right = executor.retention().unwrap()[1].pin(6);
1669 assert!(matches!(
1670 executor.step(),
1671 Err(MergeError::SourceRetained { pins: 1, .. })
1672 ));
1673 assert!(executor.retention().unwrap()[1].unpin(pin_right));
1674 assert_eq!(executor.step().unwrap(), MergePhase::SourcesRetired);
1675 assert_merge_completed(&fixture);
1676 executor.run().unwrap();
1677 assert!(MergeExecutor::resume(
1678 fixture.source_layouts.clone(),
1679 fixture.meta.clone(),
1680 fixture.keyspaces.clone(),
1681 fixture.sink.clone(),
1682 )
1683 .unwrap()
1684 .is_none());
1685 }
1686
1687 #[test]
1690 fn merge_resumes_after_a_crash_at_every_durable_boundary() {
1691 let _lock = EXECUTOR_TEST_LOCK.lock().unwrap();
1692 let hooks = [
1693 "tablet.merge.phase.1",
1694 "tablet.merge.phase.2",
1695 "tablet.merge.phase.3",
1696 "tablet.merge.phase.4",
1697 "tablet.merge.phase.5",
1698 "tablet.merge.phase.6",
1699 "tablet.merge.phase.7",
1700 "tablet.merge.before",
1701 "tablet.merge.after",
1702 ];
1703 for hook in hooks {
1704 let fixture = merge_fixture();
1705 let mut executor = begin_executor(&fixture);
1706 {
1707 let _guard =
1708 mongreldb_fault::ScopedGuard::limited(hook, mongreldb_fault::Action::Fail, 1);
1709 assert!(
1710 matches!(executor.run(), Err(MergeError::Fault(_))),
1711 "hook {hook} did not fire"
1712 );
1713 }
1714 drop(executor);
1715 let resumed = MergeExecutor::resume(
1716 fixture.source_layouts.clone(),
1717 fixture.meta.clone(),
1718 fixture.keyspaces.clone(),
1719 fixture.sink.clone(),
1720 )
1721 .unwrap();
1722 if hook == "tablet.merge.phase.7" {
1723 assert!(resumed.is_none(), "hook {hook}");
1724 } else {
1725 resumed
1726 .unwrap()
1727 .run()
1728 .unwrap_or_else(|error| panic!("resume after {hook} failed: {error}"));
1729 }
1730 assert_merge_completed(&fixture);
1731 }
1732 }
1733
1734 #[test]
1735 fn merge_build_fails_closed_on_out_of_range_keys() {
1736 let _lock = EXECUTOR_TEST_LOCK.lock().unwrap();
1737 for index in 0..2 {
1741 let fixture = merge_fixture();
1742 fixture.keyspaces[index].insert(text_key("zz-outside"), ts(120), b"rogue".to_vec());
1743 let mut executor = begin_executor(&fixture);
1744 assert!(
1745 matches!(
1746 executor.run_until(MergePhase::ReplacementBuilt),
1747 Err(MergeError::KeyOutsideSource(_))
1748 ),
1749 "source {index} accepted an out-of-range key"
1750 );
1751 }
1752 }
1753
1754 #[test]
1755 fn merge_progress_record_fails_closed_on_corruption() {
1756 let _lock = EXECUTOR_TEST_LOCK.lock().unwrap();
1757 let fixture = merge_fixture();
1758 let executor = begin_executor(&fixture);
1759 drop(executor);
1760 let path = fixture.source_layouts[0]
1761 .tablet_dir()
1762 .join(MERGE_PROGRESS_FILENAME);
1763 std::fs::write(&path, b"{ not json").unwrap();
1764 assert!(matches!(
1765 MergeExecutor::resume(
1766 fixture.source_layouts.clone(),
1767 fixture.meta.clone(),
1768 fixture.keyspaces.clone(),
1769 fixture.sink.clone(),
1770 ),
1771 Err(MergeError::Tablet(TabletError::Metadata(
1772 ClusterError::CorruptMetadata { .. }
1773 )))
1774 ));
1775 assert!(MergeExecutor::resume(
1777 [
1778 fixture.source_layouts[1].clone(),
1779 fixture.source_layouts[1].clone(),
1780 ],
1781 fixture.meta.clone(),
1782 fixture.keyspaces.clone(),
1783 fixture.sink.clone(),
1784 )
1785 .unwrap()
1786 .is_none());
1787 }
1788}