1use std::collections::{HashMap, HashSet};
18use std::time::Duration;
19
20use chrono::Utc;
21use common_telemetry::warn;
22use serde::{Deserialize, Serialize};
23use snafu::{OptionExt, ResultExt};
24use store_api::ManifestVersion;
25use store_api::metadata::RegionMetadataRef;
26use store_api::storage::{FileId, IndexVersion, RegionId, SequenceNumber};
27use strum::Display;
28
29use crate::error::{RegionMetadataNotFoundSnafu, Result, SerdeJsonSnafu, Utf8Snafu};
30use crate::manifest::manager::RemoveFileOptions;
31use crate::region::ManifestStats;
32use crate::sst::FormatType;
33use crate::sst::file::FileMeta;
34use crate::wal::EntryId;
35
36#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq, Display)]
38pub enum RegionMetaAction {
39 Change(RegionChange),
41 PartitionExprChange(RegionPartitionExprChange),
43 Edit(RegionEdit),
45 Remove(RegionRemove),
47 Truncate(RegionTruncate),
49}
50
51#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
52pub struct RegionPartitionExprChange {
53 pub partition_expr: Option<String>,
55}
56
57impl RegionMetaAction {
58 pub fn is_change(&self) -> bool {
60 matches!(self, RegionMetaAction::Change(_))
61 }
62
63 pub fn is_edit(&self) -> bool {
65 matches!(self, RegionMetaAction::Edit(_))
66 }
67
68 pub fn is_partition_expr_change(&self) -> bool {
70 matches!(self, RegionMetaAction::PartitionExprChange(_))
71 }
72}
73
74#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
75pub struct RegionChange {
76 pub metadata: RegionMetadataRef,
78 #[serde(default)]
80 pub sst_format: FormatType,
81 #[serde(default)]
83 pub append_mode: Option<bool>,
84}
85
86#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
87pub struct RegionEdit {
88 pub files_to_add: Vec<FileMeta>,
89 pub files_to_remove: Vec<FileMeta>,
90 #[serde(default)]
92 pub timestamp_ms: Option<i64>,
93 #[serde(with = "humantime_serde")]
94 pub compaction_time_window: Option<Duration>,
95 pub flushed_entry_id: Option<EntryId>,
96 pub flushed_sequence: Option<SequenceNumber>,
97 pub committed_sequence: Option<SequenceNumber>,
98}
99
100#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
101pub struct RegionRemove {
102 pub region_id: RegionId,
103}
104
105#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
108pub struct RegionTruncate {
109 pub region_id: RegionId,
110 #[serde(flatten)]
111 pub kind: TruncateKind,
112 #[serde(default)]
114 pub timestamp_ms: Option<i64>,
115}
116
117#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
119#[serde(untagged)]
120pub enum TruncateKind {
121 All {
123 truncated_entry_id: EntryId,
125 truncated_sequence: SequenceNumber,
127 },
128 Partial { files_to_remove: Vec<FileMeta> },
130}
131
132#[derive(Serialize, Deserialize, Clone, Debug)]
134#[cfg_attr(test, derive(Eq))]
135pub struct RegionManifest {
136 pub metadata: RegionMetadataRef,
138 pub files: HashMap<FileId, FileMeta>,
140 #[serde(default)]
150 pub removed_files: RemovedFilesRecord,
151 pub flushed_entry_id: EntryId,
153 pub flushed_sequence: SequenceNumber,
155 pub committed_sequence: Option<SequenceNumber>,
156 pub manifest_version: ManifestVersion,
158 pub truncated_entry_id: Option<EntryId>,
160 #[serde(with = "humantime_serde")]
162 pub compaction_time_window: Option<Duration>,
163 #[serde(default)]
165 pub sst_format: FormatType,
166 #[serde(default)]
168 pub append_mode: Option<bool>,
169}
170
171#[cfg(test)]
172impl PartialEq for RegionManifest {
173 fn eq(&self, other: &Self) -> bool {
174 self.metadata == other.metadata
175 && self.files == other.files
176 && self.flushed_entry_id == other.flushed_entry_id
177 && self.flushed_sequence == other.flushed_sequence
178 && self.manifest_version == other.manifest_version
179 && self.truncated_entry_id == other.truncated_entry_id
180 && self.compaction_time_window == other.compaction_time_window
181 && self.committed_sequence == other.committed_sequence
182 }
183}
184
185#[derive(Debug, Default)]
186pub struct RegionManifestBuilder {
187 metadata: Option<RegionMetadataRef>,
188 files: HashMap<FileId, FileMeta>,
189 pub removed_files: RemovedFilesRecord,
190 flushed_entry_id: EntryId,
191 flushed_sequence: SequenceNumber,
192 manifest_version: ManifestVersion,
193 truncated_entry_id: Option<EntryId>,
194 compaction_time_window: Option<Duration>,
195 committed_sequence: Option<SequenceNumber>,
196 sst_format: FormatType,
197 append_mode: Option<bool>,
198}
199
200impl RegionManifestBuilder {
201 fn removed_file(&self, file: &FileMeta) -> RemovedFile {
202 let index_version = self
203 .files
204 .get(&file.file_id)
205 .and_then(FileMeta::index_version)
206 .or_else(|| file.index_version());
207 RemovedFile::File(file.file_id, index_version)
208 }
209
210 pub fn with_checkpoint(checkpoint: Option<RegionManifest>) -> Self {
212 if let Some(s) = checkpoint {
213 Self {
214 metadata: Some(s.metadata),
215 files: s.files,
216 removed_files: s.removed_files,
217 flushed_entry_id: s.flushed_entry_id,
218 manifest_version: s.manifest_version,
219 flushed_sequence: s.flushed_sequence,
220 truncated_entry_id: s.truncated_entry_id,
221 compaction_time_window: s.compaction_time_window,
222 committed_sequence: s.committed_sequence,
223 sst_format: s.sst_format,
224 append_mode: s.append_mode,
225 }
226 } else {
227 Default::default()
228 }
229 }
230
231 pub fn apply_change(&mut self, manifest_version: ManifestVersion, change: RegionChange) {
232 self.metadata = Some(change.metadata);
233 self.manifest_version = manifest_version;
234 self.sst_format = change.sst_format;
235 self.append_mode = change.append_mode.or(self.append_mode);
237 }
238
239 pub fn apply_partition_expr_change(
245 &mut self,
246 manifest_version: ManifestVersion,
247 change: RegionPartitionExprChange,
248 ) {
249 if let Some(metadata) = &self.metadata {
250 let mut metadata = metadata.as_ref().clone();
251 metadata.set_partition_expr(change.partition_expr);
252 self.metadata = Some(metadata.into());
253 self.manifest_version = manifest_version;
254 } else {
255 warn!(
256 "metadata is not set in region manifest builder, ignore partition expr change: {:?}",
257 change
258 );
259 }
260 }
261
262 pub fn apply_edit(&mut self, manifest_version: ManifestVersion, edit: RegionEdit) {
263 self.manifest_version = manifest_version;
264
265 let mut removed_files = vec![];
266 for file in edit.files_to_add {
267 if let Some(old_file) = self.files.insert(file.file_id, file.clone())
268 && let Some(old_index) = old_file.index_version()
269 && !old_file.is_index_up_to_date(&file)
270 {
271 removed_files.push(RemovedFile::Index(old_file.file_id, old_index));
273 }
274 }
275 removed_files.extend(
276 edit.files_to_remove
277 .iter()
278 .map(|file| self.removed_file(file)),
279 );
280 let at = edit
281 .timestamp_ms
282 .unwrap_or_else(|| Utc::now().timestamp_millis());
283 self.removed_files.add_removed_files(removed_files, at);
284
285 for file in edit.files_to_remove {
286 self.files.remove(&file.file_id);
287 }
288 if let Some(flushed_entry_id) = edit.flushed_entry_id {
289 self.flushed_entry_id = self.flushed_entry_id.max(flushed_entry_id);
290 }
291 if let Some(flushed_sequence) = edit.flushed_sequence {
292 self.flushed_sequence = self.flushed_sequence.max(flushed_sequence);
293 }
294
295 if let Some(committed_sequence) = edit.committed_sequence {
296 self.committed_sequence = Some(
297 self.committed_sequence
298 .map_or(committed_sequence, |exist| exist.max(committed_sequence)),
299 );
300 }
301 if let Some(window) = edit.compaction_time_window {
302 self.compaction_time_window = Some(window);
303 }
304 }
305
306 pub fn apply_truncate(&mut self, manifest_version: ManifestVersion, truncate: RegionTruncate) {
307 self.manifest_version = manifest_version;
308 match truncate.kind {
309 TruncateKind::All {
310 truncated_entry_id,
311 truncated_sequence,
312 } => {
313 self.flushed_entry_id = truncated_entry_id;
314 self.flushed_sequence = truncated_sequence;
315 self.truncated_entry_id = Some(truncated_entry_id);
316 self.removed_files.add_removed_files(
317 self.files
318 .values()
319 .map(|f| RemovedFile::File(f.file_id, f.index_version()))
320 .collect(),
321 truncate
322 .timestamp_ms
323 .unwrap_or_else(|| Utc::now().timestamp_millis()),
324 );
325 self.files.clear();
326 }
327 TruncateKind::Partial { files_to_remove } => {
328 self.removed_files.add_removed_files(
334 files_to_remove
335 .iter()
336 .map(|file| self.removed_file(file))
337 .collect(),
338 truncate
339 .timestamp_ms
340 .unwrap_or_else(|| Utc::now().timestamp_millis()),
341 );
342 for file in files_to_remove {
343 self.files.remove(&file.file_id);
344 }
345 }
346 }
347 }
348
349 pub fn files(&self) -> &HashMap<FileId, FileMeta> {
350 &self.files
351 }
352
353 pub fn contains_metadata(&self) -> bool {
355 self.metadata.is_some()
356 }
357
358 pub fn try_build(self) -> Result<RegionManifest> {
359 let metadata = self.metadata.context(RegionMetadataNotFoundSnafu)?;
360 Ok(RegionManifest {
361 metadata,
362 files: self.files,
363 removed_files: self.removed_files,
364 flushed_entry_id: self.flushed_entry_id,
365 flushed_sequence: self.flushed_sequence,
366 committed_sequence: self.committed_sequence,
367 manifest_version: self.manifest_version,
368 truncated_entry_id: self.truncated_entry_id,
369 compaction_time_window: self.compaction_time_window,
370 sst_format: self.sst_format,
371 append_mode: self.append_mode,
372 })
373 }
374}
375
376#[derive(Serialize, Deserialize, Clone, Debug, Default, PartialEq, Eq)]
380pub struct RemovedFilesRecord {
381 pub removed_files: Vec<RemovedFiles>,
384}
385
386impl RemovedFilesRecord {
387 pub fn clear_deleted_files(&mut self, deleted_files: Vec<RemovedFile>) {
389 let deleted_file_set: HashSet<_> = HashSet::from_iter(deleted_files);
390 for files in self.removed_files.iter_mut() {
391 files
392 .files
393 .retain(|removed| !deleted_file_set.contains(removed));
394 }
395
396 self.removed_files.retain(|fs| !fs.files.is_empty());
397 }
398
399 pub fn update_file_removed_cnt_to_stats(&self, stats: &ManifestStats) {
400 let cnt = self
401 .removed_files
402 .iter()
403 .map(|r| r.files.len() as u64)
404 .sum();
405 stats
406 .file_removed_cnt
407 .store(cnt, std::sync::atomic::Ordering::Relaxed);
408 }
409}
410
411#[derive(Serialize, Deserialize, Clone, Debug, Default, PartialEq, Eq)]
412pub struct RemovedFiles {
413 pub removed_at: i64,
416 #[serde(default)]
418 pub files: HashSet<RemovedFile>,
419}
420
421#[derive(Serialize, Hash, Clone, Debug, PartialEq, Eq)]
423pub enum RemovedFile {
424 File(FileId, Option<IndexVersion>),
425 Index(FileId, IndexVersion),
426}
427
428impl<'de> Deserialize<'de> for RemovedFile {
432 fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
433 where
434 D: serde::Deserializer<'de>,
435 {
436 #[derive(Deserialize)]
437 #[serde(untagged)]
438 enum CompatRemovedFile {
439 Enum(RemovedFileEnum),
440 FileId(FileId),
441 }
442
443 #[derive(Deserialize)]
444 enum RemovedFileEnum {
445 File(FileId, Option<IndexVersion>),
446 Index(FileId, IndexVersion),
447 }
448
449 let compat = CompatRemovedFile::deserialize(deserializer)?;
450 match compat {
451 CompatRemovedFile::FileId(file_id) => Ok(RemovedFile::File(file_id, None)),
452 CompatRemovedFile::Enum(e) => match e {
453 RemovedFileEnum::File(file_id, version) => Ok(RemovedFile::File(file_id, version)),
454 RemovedFileEnum::Index(file_id, version) => {
455 Ok(RemovedFile::Index(file_id, version))
456 }
457 },
458 }
459 }
460}
461
462impl RemovedFile {
463 pub fn file_id(&self) -> FileId {
464 match self {
465 RemovedFile::File(file_id, _) => *file_id,
466 RemovedFile::Index(file_id, _) => *file_id,
467 }
468 }
469
470 pub fn index_version(&self) -> Option<IndexVersion> {
471 match self {
472 RemovedFile::File(_, index_version) => *index_version,
473 RemovedFile::Index(_, index_version) => Some(*index_version),
474 }
475 }
476}
477
478impl RemovedFilesRecord {
479 pub fn add_removed_files(&mut self, removed: Vec<RemovedFile>, at: i64) {
481 if removed.is_empty() {
482 return;
483 }
484 let files = removed.into_iter().collect();
485 self.removed_files.push(RemovedFiles {
486 removed_at: at,
487 files,
488 });
489 }
490
491 pub fn evict_old_removed_files(&mut self, opt: &RemoveFileOptions) -> Result<()> {
492 if !opt.enable_gc {
493 self.removed_files.clear();
495 return Ok(());
496 }
497
498 Ok(())
501 }
502}
503
504#[derive(Serialize, Deserialize, Debug, Clone)]
506#[cfg_attr(test, derive(PartialEq, Eq))]
507pub struct RegionCheckpoint {
508 pub last_version: ManifestVersion,
510 pub compacted_actions: usize,
512 pub checkpoint: Option<RegionManifest>,
514}
515
516impl RegionCheckpoint {
517 pub fn last_version(&self) -> ManifestVersion {
518 self.last_version
519 }
520
521 pub fn encode(&self) -> Result<Vec<u8>> {
522 let json = serde_json::to_string(&self).context(SerdeJsonSnafu)?;
523
524 Ok(json.into_bytes())
525 }
526
527 pub fn decode(bytes: &[u8]) -> Result<Self> {
528 let data = std::str::from_utf8(bytes).context(Utf8Snafu)?;
529
530 serde_json::from_str(data).context(SerdeJsonSnafu)
531 }
532}
533
534#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
535pub struct RegionMetaActionList {
536 pub actions: Vec<RegionMetaAction>,
537}
538
539impl RegionMetaActionList {
540 pub fn with_action(action: RegionMetaAction) -> Self {
541 Self {
542 actions: vec![action],
543 }
544 }
545
546 pub fn new(actions: Vec<RegionMetaAction>) -> Self {
547 Self { actions }
548 }
549
550 pub fn split_region_change_and_edit(
552 self,
553 ) -> (
554 Option<RegionPartitionExprChange>,
555 Option<RegionChange>,
556 RegionEdit,
557 ) {
558 let mut edit = RegionEdit {
559 files_to_add: Vec::new(),
560 files_to_remove: Vec::new(),
561 timestamp_ms: None,
562 compaction_time_window: None,
563 flushed_entry_id: None,
564 flushed_sequence: None,
565 committed_sequence: None,
566 };
567 let mut partition_expr_change = None;
568 let mut region_change = None;
569 for action in self.actions {
570 match action {
571 RegionMetaAction::PartitionExprChange(change) => {
572 partition_expr_change = Some(change);
573 }
574 RegionMetaAction::Change(change) => {
575 region_change = Some(change);
576 }
577 RegionMetaAction::Edit(region_edit) => {
578 edit.files_to_add.extend(region_edit.files_to_add);
580 edit.files_to_remove.extend(region_edit.files_to_remove);
581 if let Some(eid) = region_edit.flushed_entry_id {
583 edit.flushed_entry_id =
584 Some(edit.flushed_entry_id.map_or(eid, |v| v.max(eid)));
585 }
586 if let Some(seq) = region_edit.flushed_sequence {
587 edit.flushed_sequence =
588 Some(edit.flushed_sequence.map_or(seq, |v| v.max(seq)));
589 }
590 if let Some(seq) = region_edit.committed_sequence {
591 edit.committed_sequence =
592 Some(edit.committed_sequence.map_or(seq, |v| v.max(seq)));
593 }
594 if region_edit.compaction_time_window.is_some() {
596 edit.compaction_time_window = region_edit.compaction_time_window;
597 }
598 }
599 _ => {}
600 }
601 }
602
603 (partition_expr_change, region_change, edit)
604 }
605}
606
607impl RegionMetaActionList {
608 pub fn encode(&self) -> Result<Vec<u8>> {
610 let json = serde_json::to_string(&self).context(SerdeJsonSnafu)?;
611
612 Ok(json.into_bytes())
613 }
614
615 pub fn decode(bytes: &[u8]) -> Result<Self> {
616 let data = std::str::from_utf8(bytes).context(Utf8Snafu)?;
617
618 serde_json::from_str(data).context(SerdeJsonSnafu)
619 }
620}
621
622#[cfg(test)]
623mod tests {
624
625 use common_time::Timestamp;
626
627 use super::*;
628
629 #[test]
633 fn test_region_action_compatibility() {
634 let region_edit = r#"{
635 "flushed_entry_id":null,
636 "compaction_time_window":null,
637 "files_to_add":[
638 {"region_id":4402341478400,"file_id":"4b220a70-2b03-4641-9687-b65d94641208","time_range":[{"value":1451609210000,"unit":"Millisecond"},{"value":1451609520000,"unit":"Millisecond"}],"level":1,"file_size":100}
639 ],
640 "files_to_remove":[
641 {"region_id":4402341478400,"file_id":"34b6ebb9-b8a5-4a4b-b744-56f67defad02","time_range":[{"value":1451609210000,"unit":"Millisecond"},{"value":1451609520000,"unit":"Millisecond"}],"level":0,"file_size":100}
642 ]
643 }"#;
644 let _ = serde_json::from_str::<RegionEdit>(region_edit).unwrap();
645
646 let region_edit = r#"{
647 "flushed_entry_id":10,
648 "flushed_sequence":10,
649 "compaction_time_window":null,
650 "files_to_add":[
651 {"region_id":4402341478400,"file_id":"4b220a70-2b03-4641-9687-b65d94641208","time_range":[{"value":1451609210000,"unit":"Millisecond"},{"value":1451609520000,"unit":"Millisecond"}],"level":1,"file_size":100}
652 ],
653 "files_to_remove":[
654 {"region_id":4402341478400,"file_id":"34b6ebb9-b8a5-4a4b-b744-56f67defad02","time_range":[{"value":1451609210000,"unit":"Millisecond"},{"value":1451609520000,"unit":"Millisecond"}],"level":0,"file_size":100}
655 ]
656 }"#;
657 let _ = serde_json::from_str::<RegionEdit>(region_edit).unwrap();
658
659 let region_change = r#" {
661 "metadata":{
662 "column_metadatas":[
663 {"column_schema":{"name":"a","data_type":{"Int64":{}},"is_nullable":false,"is_time_index":false,"default_constraint":null,"metadata":{}},"semantic_type":"Tag","column_id":1},{"column_schema":{"name":"b","data_type":{"Float64":{}},"is_nullable":false,"is_time_index":false,"default_constraint":null,"metadata":{}},"semantic_type":"Field","column_id":2},{"column_schema":{"name":"c","data_type":{"Timestamp":{"Millisecond":null}},"is_nullable":false,"is_time_index":false,"default_constraint":null,"metadata":{}},"semantic_type":"Timestamp","column_id":3}
664 ],
665 "primary_key":[1],
666 "region_id":5299989648942,
667 "schema_version":0
668 }
669 }"#;
670 let _ = serde_json::from_str::<RegionChange>(region_change).unwrap();
671
672 let region_remove = r#"{"region_id":42}"#;
673 let _ = serde_json::from_str::<RegionRemove>(region_remove).unwrap();
674
675 let region_partition_expr_change = r#"{
676 "partition_expr": "{\"expr\":\"x < 100\"}"
677 }"#;
678 let _ = serde_json::from_str::<RegionPartitionExprChange>(region_partition_expr_change)
679 .unwrap();
680 }
681
682 #[test]
683 fn test_region_manifest_compatibility() {
684 let region_manifest_json = r#"{
686 "metadata": {
687 "column_metadatas": [
688 {
689 "column_schema": {
690 "name": "a",
691 "data_type": {"Int64": {}},
692 "is_nullable": false,
693 "is_time_index": false,
694 "default_constraint": null,
695 "metadata": {}
696 },
697 "semantic_type": "Tag",
698 "column_id": 1
699 },
700 {
701 "column_schema": {
702 "name": "b",
703 "data_type": {"Float64": {}},
704 "is_nullable": false,
705 "is_time_index": false,
706 "default_constraint": null,
707 "metadata": {}
708 },
709 "semantic_type": "Field",
710 "column_id": 2
711 },
712 {
713 "column_schema": {
714 "name": "c",
715 "data_type": {"Timestamp": {"Millisecond": null}},
716 "is_nullable": false,
717 "is_time_index": false,
718 "default_constraint": null,
719 "metadata": {}
720 },
721 "semantic_type": "Timestamp",
722 "column_id": 3
723 }
724 ],
725 "primary_key": [1],
726 "region_id": 4402341478400,
727 "schema_version": 0
728 },
729 "files": {
730 "4b220a70-2b03-4641-9687-b65d94641208": {
731 "region_id": 4402341478400,
732 "file_id": "4b220a70-2b03-4641-9687-b65d94641208",
733 "time_range": [
734 {"value": 1451609210000, "unit": "Millisecond"},
735 {"value": 1451609520000, "unit": "Millisecond"}
736 ],
737 "level": 1,
738 "file_size": 100
739 },
740 "34b6ebb9-b8a5-4a4b-b744-56f67defad02": {
741 "region_id": 4402341478400,
742 "file_id": "34b6ebb9-b8a5-4a4b-b744-56f67defad02",
743 "time_range": [
744 {"value": 1451609210000, "unit": "Millisecond"},
745 {"value": 1451609520000, "unit": "Millisecond"}
746 ],
747 "level": 0,
748 "file_size": 100
749 }
750 },
751 "flushed_entry_id": 10,
752 "flushed_sequence": 20,
753 "manifest_version": 1,
754 "truncated_entry_id": null,
755 "compaction_time_window": null
756 }"#;
757
758 let manifest = serde_json::from_str::<RegionManifest>(region_manifest_json).unwrap();
759
760 assert_eq!(manifest.files.len(), 2);
762 assert_eq!(manifest.flushed_entry_id, 10);
763 assert_eq!(manifest.flushed_sequence, 20);
764 assert_eq!(manifest.manifest_version, 1);
765
766 let mut file_ids: Vec<String> = manifest.files.keys().map(|id| id.to_string()).collect();
768 file_ids.sort_unstable();
769 assert_eq!(
770 file_ids,
771 vec![
772 "34b6ebb9-b8a5-4a4b-b744-56f67defad02",
773 "4b220a70-2b03-4641-9687-b65d94641208",
774 ]
775 );
776
777 let serialized_manifest = serde_json::to_string(&manifest).unwrap();
779 let deserialized_manifest: RegionManifest =
780 serde_json::from_str(&serialized_manifest).unwrap();
781 assert_eq!(manifest, deserialized_manifest);
782 assert_ne!(serialized_manifest, region_manifest_json);
783 }
784
785 #[test]
786 fn test_region_truncate_compat() {
787 let region_truncate_json = r#"{
789 "region_id": 4402341478400,
790 "truncated_entry_id": 10,
791 "truncated_sequence": 20
792 }"#;
793
794 let truncate_v1: RegionTruncate = serde_json::from_str(region_truncate_json).unwrap();
795 assert_eq!(truncate_v1.region_id, 4402341478400);
796 assert_eq!(
797 truncate_v1.kind,
798 TruncateKind::All {
799 truncated_entry_id: 10,
800 truncated_sequence: 20,
801 }
802 );
803
804 let region_truncate_v2_json = r#"{
806 "region_id": 4402341478400,
807 "files_to_remove": [
808 {
809 "region_id": 4402341478400,
810 "file_id": "4b220a70-2b03-4641-9687-b65d94641208",
811 "time_range": [
812 {
813 "value": 1451609210000,
814 "unit": "Millisecond"
815 },
816 {
817 "value": 1451609520000,
818 "unit": "Millisecond"
819 }
820 ],
821 "level": 1,
822 "file_size": 100
823 }
824 ]
825}"#;
826
827 let truncate_v2: RegionTruncate = serde_json::from_str(region_truncate_v2_json).unwrap();
828 assert_eq!(truncate_v2.region_id, 4402341478400);
829 assert_eq!(
830 truncate_v2.kind,
831 TruncateKind::Partial {
832 files_to_remove: vec![FileMeta {
833 region_id: RegionId::from_u64(4402341478400),
834 file_id: FileId::parse_str("4b220a70-2b03-4641-9687-b65d94641208").unwrap(),
835 time_range: (
836 Timestamp::new_millisecond(1451609210000),
837 Timestamp::new_millisecond(1451609520000)
838 ),
839 level: 1,
840 file_size: 100,
841 ..Default::default()
842 }]
843 }
844 );
845 }
846
847 #[test]
848 fn test_region_manifest_removed_files() {
849 let region_metadata = r#"{
850 "column_metadatas": [
851 {
852 "column_schema": {
853 "name": "a",
854 "data_type": {"Int64": {}},
855 "is_nullable": false,
856 "is_time_index": false,
857 "default_constraint": null,
858 "metadata": {}
859 },
860 "semantic_type": "Tag",
861 "column_id": 1
862 },
863 {
864 "column_schema": {
865 "name": "b",
866 "data_type": {"Float64": {}},
867 "is_nullable": false,
868 "is_time_index": false,
869 "default_constraint": null,
870 "metadata": {}
871 },
872 "semantic_type": "Field",
873 "column_id": 2
874 },
875 {
876 "column_schema": {
877 "name": "c",
878 "data_type": {"Timestamp": {"Millisecond": null}},
879 "is_nullable": false,
880 "is_time_index": false,
881 "default_constraint": null,
882 "metadata": {}
883 },
884 "semantic_type": "Timestamp",
885 "column_id": 3
886 }
887 ],
888 "primary_key": [1],
889 "region_id": 4402341478400,
890 "schema_version": 0
891 }"#;
892
893 let metadata: RegionMetadataRef =
894 serde_json::from_str(region_metadata).expect("Failed to parse region metadata");
895 let manifest = RegionManifest {
896 metadata: metadata.clone(),
897 files: HashMap::new(),
898 flushed_entry_id: 0,
899 flushed_sequence: 0,
900 committed_sequence: None,
901 manifest_version: 0,
902 truncated_entry_id: None,
903 compaction_time_window: None,
904 removed_files: RemovedFilesRecord {
905 removed_files: vec![RemovedFiles {
906 removed_at: 0,
907 files: HashSet::from([RemovedFile::File(
908 FileId::parse_str("4b220a70-2b03-4641-9687-b65d94641208").unwrap(),
909 None,
910 )]),
911 }],
912 },
913 sst_format: FormatType::PrimaryKey,
914 append_mode: None,
915 };
916
917 let json = serde_json::to_string(&manifest).unwrap();
918 let new: RegionManifest = serde_json::from_str(&json).unwrap();
919
920 assert_eq!(manifest, new);
921 }
922
923 #[test]
924 fn test_remove_tracks_current_manifest_index_version() {
925 let file_id = FileId::random();
926 let current = FileMeta {
927 region_id: RegionId::new(1, 1),
928 file_id,
929 available_indexes: smallvec::smallvec![crate::sst::file::IndexType::InvertedIndex],
930 index_file_size: 1024,
931 index_version: 2,
932 ..Default::default()
933 };
934 let mut stale = current.clone();
935 stale.available_indexes.clear();
936 stale.index_file_size = 0;
937 stale.index_version = 0;
938
939 let mut builder = RegionManifestBuilder::default();
940 builder.apply_edit(
941 1,
942 RegionEdit {
943 files_to_add: vec![current.clone()],
944 files_to_remove: Vec::new(),
945 timestamp_ms: None,
946 compaction_time_window: None,
947 flushed_entry_id: None,
948 flushed_sequence: None,
949 committed_sequence: None,
950 },
951 );
952 builder.apply_edit(
953 2,
954 RegionEdit {
955 files_to_add: Vec::new(),
956 files_to_remove: vec![stale.clone()],
957 timestamp_ms: Some(42),
958 compaction_time_window: None,
959 flushed_entry_id: None,
960 flushed_sequence: None,
961 committed_sequence: None,
962 },
963 );
964
965 assert_eq!(
966 builder.removed_files.removed_files[0].files,
967 HashSet::from([RemovedFile::File(file_id, Some(2))])
968 );
969
970 let mut builder = RegionManifestBuilder::default();
971 builder.apply_edit(
972 1,
973 RegionEdit {
974 files_to_add: vec![current],
975 files_to_remove: Vec::new(),
976 timestamp_ms: None,
977 compaction_time_window: None,
978 flushed_entry_id: None,
979 flushed_sequence: None,
980 committed_sequence: None,
981 },
982 );
983 builder.apply_truncate(
984 2,
985 RegionTruncate {
986 region_id: RegionId::new(1, 1),
987 kind: TruncateKind::Partial {
988 files_to_remove: vec![stale],
989 },
990 timestamp_ms: Some(42),
991 },
992 );
993
994 assert_eq!(
995 builder.removed_files.removed_files[0].files,
996 HashSet::from([RemovedFile::File(file_id, Some(2))])
997 );
998 }
999
1000 #[test]
1002 fn test_old_region_manifest_compat() {
1003 #[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
1004 pub struct RegionManifestV1 {
1005 pub metadata: RegionMetadataRef,
1007 pub files: HashMap<FileId, FileMeta>,
1009 pub flushed_entry_id: EntryId,
1011 pub flushed_sequence: SequenceNumber,
1013 pub manifest_version: ManifestVersion,
1015 pub truncated_entry_id: Option<EntryId>,
1017 #[serde(with = "humantime_serde")]
1019 pub compaction_time_window: Option<Duration>,
1020 }
1021
1022 let region_metadata = r#"{
1023 "column_metadatas": [
1024 {
1025 "column_schema": {
1026 "name": "a",
1027 "data_type": {"Int64": {}},
1028 "is_nullable": false,
1029 "is_time_index": false,
1030 "default_constraint": null,
1031 "metadata": {}
1032 },
1033 "semantic_type": "Tag",
1034 "column_id": 1
1035 },
1036 {
1037 "column_schema": {
1038 "name": "b",
1039 "data_type": {"Float64": {}},
1040 "is_nullable": false,
1041 "is_time_index": false,
1042 "default_constraint": null,
1043 "metadata": {}
1044 },
1045 "semantic_type": "Field",
1046 "column_id": 2
1047 },
1048 {
1049 "column_schema": {
1050 "name": "c",
1051 "data_type": {"Timestamp": {"Millisecond": null}},
1052 "is_nullable": false,
1053 "is_time_index": false,
1054 "default_constraint": null,
1055 "metadata": {}
1056 },
1057 "semantic_type": "Timestamp",
1058 "column_id": 3
1059 }
1060 ],
1061 "primary_key": [1],
1062 "region_id": 4402341478400,
1063 "schema_version": 0
1064 }"#;
1065
1066 let metadata: RegionMetadataRef =
1067 serde_json::from_str(region_metadata).expect("Failed to parse region metadata");
1068
1069 let v1 = RegionManifestV1 {
1071 metadata: metadata.clone(),
1072 files: HashMap::new(),
1073 flushed_entry_id: 0,
1074 flushed_sequence: 0,
1075 manifest_version: 0,
1076 truncated_entry_id: None,
1077 compaction_time_window: None,
1078 };
1079 let json = serde_json::to_string(&v1).unwrap();
1080 let new_from_old: RegionManifest = serde_json::from_str(&json).unwrap();
1081 assert_eq!(
1082 new_from_old,
1083 RegionManifest {
1084 metadata: metadata.clone(),
1085 files: HashMap::new(),
1086 removed_files: Default::default(),
1087 flushed_entry_id: 0,
1088 flushed_sequence: 0,
1089 committed_sequence: None,
1090 manifest_version: 0,
1091 truncated_entry_id: None,
1092 compaction_time_window: None,
1093 sst_format: FormatType::PrimaryKey,
1094 append_mode: None,
1095 }
1096 );
1097
1098 let new_manifest = RegionManifest {
1099 metadata: metadata.clone(),
1100 files: HashMap::new(),
1101 removed_files: Default::default(),
1102 flushed_entry_id: 0,
1103 flushed_sequence: 0,
1104 committed_sequence: None,
1105 manifest_version: 0,
1106 truncated_entry_id: None,
1107 compaction_time_window: None,
1108 sst_format: FormatType::PrimaryKey,
1109 append_mode: None,
1110 };
1111 let json = serde_json::to_string(&new_manifest).unwrap();
1112 let old_from_new: RegionManifestV1 = serde_json::from_str(&json).unwrap();
1113 assert_eq!(
1114 old_from_new,
1115 RegionManifestV1 {
1116 metadata: metadata.clone(),
1117 files: HashMap::new(),
1118 flushed_entry_id: 0,
1119 flushed_sequence: 0,
1120 manifest_version: 0,
1121 truncated_entry_id: None,
1122 compaction_time_window: None,
1123 }
1124 );
1125
1126 #[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
1127 pub struct RegionEditV1 {
1128 pub files_to_add: Vec<FileMeta>,
1129 pub files_to_remove: Vec<FileMeta>,
1130 #[serde(with = "humantime_serde")]
1131 pub compaction_time_window: Option<Duration>,
1132 pub flushed_entry_id: Option<EntryId>,
1133 pub flushed_sequence: Option<SequenceNumber>,
1134 }
1135
1136 let json = serde_json::to_string(&RegionEditV1 {
1137 files_to_add: vec![],
1138 files_to_remove: vec![],
1139 compaction_time_window: None,
1140 flushed_entry_id: None,
1141 flushed_sequence: None,
1142 })
1143 .unwrap();
1144 let new_from_old: RegionEdit = serde_json::from_str(&json).unwrap();
1145 assert_eq!(
1146 RegionEdit {
1147 files_to_add: vec![],
1148 files_to_remove: vec![],
1149 timestamp_ms: None,
1150 compaction_time_window: None,
1151 flushed_entry_id: None,
1152 flushed_sequence: None,
1153 committed_sequence: None,
1154 },
1155 new_from_old
1156 );
1157
1158 let new = RegionEdit {
1160 files_to_add: vec![],
1161 files_to_remove: vec![],
1162 timestamp_ms: Some(42),
1163 compaction_time_window: None,
1164 flushed_entry_id: None,
1165 flushed_sequence: None,
1166 committed_sequence: None,
1167 };
1168
1169 let new_json = serde_json::to_string(&new).unwrap();
1170
1171 let old_from_new: RegionEditV1 = serde_json::from_str(&new_json).unwrap();
1172 assert_eq!(
1173 RegionEditV1 {
1174 files_to_add: vec![],
1175 files_to_remove: vec![],
1176 compaction_time_window: None,
1177 flushed_entry_id: None,
1178 flushed_sequence: None,
1179 },
1180 old_from_new
1181 );
1182 }
1183
1184 #[test]
1185 fn test_region_change_backward_compatibility() {
1186 let region_change_json = r#"{
1188 "metadata": {
1189 "column_metadatas": [
1190 {"column_schema":{"name":"a","data_type":{"Int64":{}},"is_nullable":false,"is_time_index":false,"default_constraint":null,"metadata":{}},"semantic_type":"Tag","column_id":1},
1191 {"column_schema":{"name":"b","data_type":{"Int64":{}},"is_nullable":false,"is_time_index":false,"default_constraint":null,"metadata":{}},"semantic_type":"Field","column_id":2},
1192 {"column_schema":{"name":"c","data_type":{"Timestamp":{"Millisecond":null}},"is_nullable":false,"is_time_index":false,"default_constraint":null,"metadata":{}},"semantic_type":"Timestamp","column_id":3}
1193 ],
1194 "primary_key": [
1195 1
1196 ],
1197 "region_id": 42,
1198 "schema_version": 0
1199 }
1200 }"#;
1201
1202 let region_change: RegionChange = serde_json::from_str(region_change_json).unwrap();
1203 assert_eq!(region_change.sst_format, FormatType::PrimaryKey);
1204
1205 let region_change = RegionChange {
1207 metadata: region_change.metadata.clone(),
1208 sst_format: FormatType::Flat,
1209 append_mode: None,
1210 };
1211
1212 let serialized = serde_json::to_string(®ion_change).unwrap();
1213 let deserialized: RegionChange = serde_json::from_str(&serialized).unwrap();
1214 assert_eq!(deserialized.sst_format, FormatType::Flat);
1215 }
1216
1217 #[test]
1218 fn test_removed_file_compatibility() {
1219 let file_id = FileId::random();
1220 let json_str = format!("\"{}\"", file_id);
1222 let removed_file: RemovedFile = serde_json::from_str(&json_str).unwrap();
1223 assert_eq!(removed_file, RemovedFile::File(file_id, None));
1224
1225 let removed_file_v2 = RemovedFile::File(file_id, Some(10));
1227 let json_v2 = serde_json::to_string(&removed_file_v2).unwrap();
1228 let deserialized_v2: RemovedFile = serde_json::from_str(&json_v2).unwrap();
1229 assert_eq!(removed_file_v2, deserialized_v2);
1230
1231 let removed_index = RemovedFile::Index(file_id, 20);
1233 let json_index = serde_json::to_string(&removed_index).unwrap();
1234 let deserialized_index: RemovedFile = serde_json::from_str(&json_index).unwrap();
1235 assert_eq!(removed_index, deserialized_index);
1236
1237 let removed_file = RemovedFile::File(file_id, None);
1239 let json = serde_json::to_string(&removed_file).unwrap();
1240 let deserialized: RemovedFile = serde_json::from_str(&json).unwrap();
1241 assert_eq!(removed_file, deserialized);
1242
1243 let json_set = format!("[\"{}\"]", file_id);
1249 let removed_files_set: HashSet<RemovedFile> = serde_json::from_str(&json_set).unwrap();
1250 assert!(removed_files_set.contains(&RemovedFile::File(file_id, None)));
1251 }
1252
1253 #[test]
1276 fn test_removed_files_backward_compatibility() {
1277 #[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
1279 struct OldRemovedFiles {
1280 pub removed_at: i64,
1281 pub file_ids: HashSet<FileId>,
1282 }
1283
1284 let mut file_ids = HashSet::new();
1286 file_ids.insert(FileId::random());
1287 file_ids.insert(FileId::random());
1288
1289 let old_removed_files = OldRemovedFiles {
1290 removed_at: 1234567890,
1291 file_ids,
1292 };
1293
1294 let old_json = serde_json::to_string(&old_removed_files).unwrap();
1296
1297 let result: Result<RemovedFiles, _> = serde_json::from_str(&old_json);
1299
1300 assert!(result.is_ok(), "{:?}", result);
1302 let removed_files = result.unwrap();
1303 assert_eq!(removed_files.removed_at, 1234567890);
1304 assert!(removed_files.files.is_empty());
1305
1306 let file_id = FileId::random();
1308 let new_json = format!(
1309 r#"{{
1310 "removed_at": 1234567890,
1311 "files": ["{}"]
1312 }}"#,
1313 file_id
1314 );
1315
1316 let result: Result<RemovedFiles, _> = serde_json::from_str(&new_json);
1317 assert!(result.is_ok());
1318 let removed_files = result.unwrap();
1319 assert_eq!(removed_files.removed_at, 1234567890);
1320 assert_eq!(removed_files.files.len(), 1);
1321 assert!(
1322 removed_files
1323 .files
1324 .contains(&RemovedFile::File(file_id, None))
1325 );
1326 }
1327}