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_builder() {
684 }
686
687 #[test]
688 fn test_encode_decode_region_checkpoint() {
689 }
691
692 #[test]
693 fn test_region_manifest_compatibility() {
694 let region_manifest_json = r#"{
696 "metadata": {
697 "column_metadatas": [
698 {
699 "column_schema": {
700 "name": "a",
701 "data_type": {"Int64": {}},
702 "is_nullable": false,
703 "is_time_index": false,
704 "default_constraint": null,
705 "metadata": {}
706 },
707 "semantic_type": "Tag",
708 "column_id": 1
709 },
710 {
711 "column_schema": {
712 "name": "b",
713 "data_type": {"Float64": {}},
714 "is_nullable": false,
715 "is_time_index": false,
716 "default_constraint": null,
717 "metadata": {}
718 },
719 "semantic_type": "Field",
720 "column_id": 2
721 },
722 {
723 "column_schema": {
724 "name": "c",
725 "data_type": {"Timestamp": {"Millisecond": null}},
726 "is_nullable": false,
727 "is_time_index": false,
728 "default_constraint": null,
729 "metadata": {}
730 },
731 "semantic_type": "Timestamp",
732 "column_id": 3
733 }
734 ],
735 "primary_key": [1],
736 "region_id": 4402341478400,
737 "schema_version": 0
738 },
739 "files": {
740 "4b220a70-2b03-4641-9687-b65d94641208": {
741 "region_id": 4402341478400,
742 "file_id": "4b220a70-2b03-4641-9687-b65d94641208",
743 "time_range": [
744 {"value": 1451609210000, "unit": "Millisecond"},
745 {"value": 1451609520000, "unit": "Millisecond"}
746 ],
747 "level": 1,
748 "file_size": 100
749 },
750 "34b6ebb9-b8a5-4a4b-b744-56f67defad02": {
751 "region_id": 4402341478400,
752 "file_id": "34b6ebb9-b8a5-4a4b-b744-56f67defad02",
753 "time_range": [
754 {"value": 1451609210000, "unit": "Millisecond"},
755 {"value": 1451609520000, "unit": "Millisecond"}
756 ],
757 "level": 0,
758 "file_size": 100
759 }
760 },
761 "flushed_entry_id": 10,
762 "flushed_sequence": 20,
763 "manifest_version": 1,
764 "truncated_entry_id": null,
765 "compaction_time_window": null
766 }"#;
767
768 let manifest = serde_json::from_str::<RegionManifest>(region_manifest_json).unwrap();
769
770 assert_eq!(manifest.files.len(), 2);
772 assert_eq!(manifest.flushed_entry_id, 10);
773 assert_eq!(manifest.flushed_sequence, 20);
774 assert_eq!(manifest.manifest_version, 1);
775
776 let mut file_ids: Vec<String> = manifest.files.keys().map(|id| id.to_string()).collect();
778 file_ids.sort_unstable();
779 assert_eq!(
780 file_ids,
781 vec![
782 "34b6ebb9-b8a5-4a4b-b744-56f67defad02",
783 "4b220a70-2b03-4641-9687-b65d94641208",
784 ]
785 );
786
787 let serialized_manifest = serde_json::to_string(&manifest).unwrap();
789 let deserialized_manifest: RegionManifest =
790 serde_json::from_str(&serialized_manifest).unwrap();
791 assert_eq!(manifest, deserialized_manifest);
792 assert_ne!(serialized_manifest, region_manifest_json);
793 }
794
795 #[test]
796 fn test_region_truncate_compat() {
797 let region_truncate_json = r#"{
799 "region_id": 4402341478400,
800 "truncated_entry_id": 10,
801 "truncated_sequence": 20
802 }"#;
803
804 let truncate_v1: RegionTruncate = serde_json::from_str(region_truncate_json).unwrap();
805 assert_eq!(truncate_v1.region_id, 4402341478400);
806 assert_eq!(
807 truncate_v1.kind,
808 TruncateKind::All {
809 truncated_entry_id: 10,
810 truncated_sequence: 20,
811 }
812 );
813
814 let region_truncate_v2_json = r#"{
816 "region_id": 4402341478400,
817 "files_to_remove": [
818 {
819 "region_id": 4402341478400,
820 "file_id": "4b220a70-2b03-4641-9687-b65d94641208",
821 "time_range": [
822 {
823 "value": 1451609210000,
824 "unit": "Millisecond"
825 },
826 {
827 "value": 1451609520000,
828 "unit": "Millisecond"
829 }
830 ],
831 "level": 1,
832 "file_size": 100
833 }
834 ]
835}"#;
836
837 let truncate_v2: RegionTruncate = serde_json::from_str(region_truncate_v2_json).unwrap();
838 assert_eq!(truncate_v2.region_id, 4402341478400);
839 assert_eq!(
840 truncate_v2.kind,
841 TruncateKind::Partial {
842 files_to_remove: vec![FileMeta {
843 region_id: RegionId::from_u64(4402341478400),
844 file_id: FileId::parse_str("4b220a70-2b03-4641-9687-b65d94641208").unwrap(),
845 time_range: (
846 Timestamp::new_millisecond(1451609210000),
847 Timestamp::new_millisecond(1451609520000)
848 ),
849 level: 1,
850 file_size: 100,
851 ..Default::default()
852 }]
853 }
854 );
855 }
856
857 #[test]
858 fn test_region_manifest_removed_files() {
859 let region_metadata = r#"{
860 "column_metadatas": [
861 {
862 "column_schema": {
863 "name": "a",
864 "data_type": {"Int64": {}},
865 "is_nullable": false,
866 "is_time_index": false,
867 "default_constraint": null,
868 "metadata": {}
869 },
870 "semantic_type": "Tag",
871 "column_id": 1
872 },
873 {
874 "column_schema": {
875 "name": "b",
876 "data_type": {"Float64": {}},
877 "is_nullable": false,
878 "is_time_index": false,
879 "default_constraint": null,
880 "metadata": {}
881 },
882 "semantic_type": "Field",
883 "column_id": 2
884 },
885 {
886 "column_schema": {
887 "name": "c",
888 "data_type": {"Timestamp": {"Millisecond": null}},
889 "is_nullable": false,
890 "is_time_index": false,
891 "default_constraint": null,
892 "metadata": {}
893 },
894 "semantic_type": "Timestamp",
895 "column_id": 3
896 }
897 ],
898 "primary_key": [1],
899 "region_id": 4402341478400,
900 "schema_version": 0
901 }"#;
902
903 let metadata: RegionMetadataRef =
904 serde_json::from_str(region_metadata).expect("Failed to parse region metadata");
905 let manifest = RegionManifest {
906 metadata: metadata.clone(),
907 files: HashMap::new(),
908 flushed_entry_id: 0,
909 flushed_sequence: 0,
910 committed_sequence: None,
911 manifest_version: 0,
912 truncated_entry_id: None,
913 compaction_time_window: None,
914 removed_files: RemovedFilesRecord {
915 removed_files: vec![RemovedFiles {
916 removed_at: 0,
917 files: HashSet::from([RemovedFile::File(
918 FileId::parse_str("4b220a70-2b03-4641-9687-b65d94641208").unwrap(),
919 None,
920 )]),
921 }],
922 },
923 sst_format: FormatType::PrimaryKey,
924 append_mode: None,
925 };
926
927 let json = serde_json::to_string(&manifest).unwrap();
928 let new: RegionManifest = serde_json::from_str(&json).unwrap();
929
930 assert_eq!(manifest, new);
931 }
932
933 #[test]
934 fn test_remove_tracks_current_manifest_index_version() {
935 let file_id = FileId::random();
936 let current = FileMeta {
937 region_id: RegionId::new(1, 1),
938 file_id,
939 available_indexes: smallvec::smallvec![crate::sst::file::IndexType::InvertedIndex],
940 index_file_size: 1024,
941 index_version: 2,
942 ..Default::default()
943 };
944 let mut stale = current.clone();
945 stale.available_indexes.clear();
946 stale.index_file_size = 0;
947 stale.index_version = 0;
948
949 let mut builder = RegionManifestBuilder::default();
950 builder.apply_edit(
951 1,
952 RegionEdit {
953 files_to_add: vec![current.clone()],
954 files_to_remove: Vec::new(),
955 timestamp_ms: None,
956 compaction_time_window: None,
957 flushed_entry_id: None,
958 flushed_sequence: None,
959 committed_sequence: None,
960 },
961 );
962 builder.apply_edit(
963 2,
964 RegionEdit {
965 files_to_add: Vec::new(),
966 files_to_remove: vec![stale.clone()],
967 timestamp_ms: Some(42),
968 compaction_time_window: None,
969 flushed_entry_id: None,
970 flushed_sequence: None,
971 committed_sequence: None,
972 },
973 );
974
975 assert_eq!(
976 builder.removed_files.removed_files[0].files,
977 HashSet::from([RemovedFile::File(file_id, Some(2))])
978 );
979
980 let mut builder = RegionManifestBuilder::default();
981 builder.apply_edit(
982 1,
983 RegionEdit {
984 files_to_add: vec![current],
985 files_to_remove: Vec::new(),
986 timestamp_ms: None,
987 compaction_time_window: None,
988 flushed_entry_id: None,
989 flushed_sequence: None,
990 committed_sequence: None,
991 },
992 );
993 builder.apply_truncate(
994 2,
995 RegionTruncate {
996 region_id: RegionId::new(1, 1),
997 kind: TruncateKind::Partial {
998 files_to_remove: vec![stale],
999 },
1000 timestamp_ms: Some(42),
1001 },
1002 );
1003
1004 assert_eq!(
1005 builder.removed_files.removed_files[0].files,
1006 HashSet::from([RemovedFile::File(file_id, Some(2))])
1007 );
1008 }
1009
1010 #[test]
1012 fn test_old_region_manifest_compat() {
1013 #[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
1014 pub struct RegionManifestV1 {
1015 pub metadata: RegionMetadataRef,
1017 pub files: HashMap<FileId, FileMeta>,
1019 pub flushed_entry_id: EntryId,
1021 pub flushed_sequence: SequenceNumber,
1023 pub manifest_version: ManifestVersion,
1025 pub truncated_entry_id: Option<EntryId>,
1027 #[serde(with = "humantime_serde")]
1029 pub compaction_time_window: Option<Duration>,
1030 }
1031
1032 let region_metadata = r#"{
1033 "column_metadatas": [
1034 {
1035 "column_schema": {
1036 "name": "a",
1037 "data_type": {"Int64": {}},
1038 "is_nullable": false,
1039 "is_time_index": false,
1040 "default_constraint": null,
1041 "metadata": {}
1042 },
1043 "semantic_type": "Tag",
1044 "column_id": 1
1045 },
1046 {
1047 "column_schema": {
1048 "name": "b",
1049 "data_type": {"Float64": {}},
1050 "is_nullable": false,
1051 "is_time_index": false,
1052 "default_constraint": null,
1053 "metadata": {}
1054 },
1055 "semantic_type": "Field",
1056 "column_id": 2
1057 },
1058 {
1059 "column_schema": {
1060 "name": "c",
1061 "data_type": {"Timestamp": {"Millisecond": null}},
1062 "is_nullable": false,
1063 "is_time_index": false,
1064 "default_constraint": null,
1065 "metadata": {}
1066 },
1067 "semantic_type": "Timestamp",
1068 "column_id": 3
1069 }
1070 ],
1071 "primary_key": [1],
1072 "region_id": 4402341478400,
1073 "schema_version": 0
1074 }"#;
1075
1076 let metadata: RegionMetadataRef =
1077 serde_json::from_str(region_metadata).expect("Failed to parse region metadata");
1078
1079 let v1 = RegionManifestV1 {
1081 metadata: metadata.clone(),
1082 files: HashMap::new(),
1083 flushed_entry_id: 0,
1084 flushed_sequence: 0,
1085 manifest_version: 0,
1086 truncated_entry_id: None,
1087 compaction_time_window: None,
1088 };
1089 let json = serde_json::to_string(&v1).unwrap();
1090 let new_from_old: RegionManifest = serde_json::from_str(&json).unwrap();
1091 assert_eq!(
1092 new_from_old,
1093 RegionManifest {
1094 metadata: metadata.clone(),
1095 files: HashMap::new(),
1096 removed_files: Default::default(),
1097 flushed_entry_id: 0,
1098 flushed_sequence: 0,
1099 committed_sequence: None,
1100 manifest_version: 0,
1101 truncated_entry_id: None,
1102 compaction_time_window: None,
1103 sst_format: FormatType::PrimaryKey,
1104 append_mode: None,
1105 }
1106 );
1107
1108 let new_manifest = RegionManifest {
1109 metadata: metadata.clone(),
1110 files: HashMap::new(),
1111 removed_files: Default::default(),
1112 flushed_entry_id: 0,
1113 flushed_sequence: 0,
1114 committed_sequence: None,
1115 manifest_version: 0,
1116 truncated_entry_id: None,
1117 compaction_time_window: None,
1118 sst_format: FormatType::PrimaryKey,
1119 append_mode: None,
1120 };
1121 let json = serde_json::to_string(&new_manifest).unwrap();
1122 let old_from_new: RegionManifestV1 = serde_json::from_str(&json).unwrap();
1123 assert_eq!(
1124 old_from_new,
1125 RegionManifestV1 {
1126 metadata: metadata.clone(),
1127 files: HashMap::new(),
1128 flushed_entry_id: 0,
1129 flushed_sequence: 0,
1130 manifest_version: 0,
1131 truncated_entry_id: None,
1132 compaction_time_window: None,
1133 }
1134 );
1135
1136 #[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
1137 pub struct RegionEditV1 {
1138 pub files_to_add: Vec<FileMeta>,
1139 pub files_to_remove: Vec<FileMeta>,
1140 #[serde(with = "humantime_serde")]
1141 pub compaction_time_window: Option<Duration>,
1142 pub flushed_entry_id: Option<EntryId>,
1143 pub flushed_sequence: Option<SequenceNumber>,
1144 }
1145
1146 let json = serde_json::to_string(&RegionEditV1 {
1147 files_to_add: vec![],
1148 files_to_remove: vec![],
1149 compaction_time_window: None,
1150 flushed_entry_id: None,
1151 flushed_sequence: None,
1152 })
1153 .unwrap();
1154 let new_from_old: RegionEdit = serde_json::from_str(&json).unwrap();
1155 assert_eq!(
1156 RegionEdit {
1157 files_to_add: vec![],
1158 files_to_remove: vec![],
1159 timestamp_ms: None,
1160 compaction_time_window: None,
1161 flushed_entry_id: None,
1162 flushed_sequence: None,
1163 committed_sequence: None,
1164 },
1165 new_from_old
1166 );
1167
1168 let new = RegionEdit {
1170 files_to_add: vec![],
1171 files_to_remove: vec![],
1172 timestamp_ms: Some(42),
1173 compaction_time_window: None,
1174 flushed_entry_id: None,
1175 flushed_sequence: None,
1176 committed_sequence: None,
1177 };
1178
1179 let new_json = serde_json::to_string(&new).unwrap();
1180
1181 let old_from_new: RegionEditV1 = serde_json::from_str(&new_json).unwrap();
1182 assert_eq!(
1183 RegionEditV1 {
1184 files_to_add: vec![],
1185 files_to_remove: vec![],
1186 compaction_time_window: None,
1187 flushed_entry_id: None,
1188 flushed_sequence: None,
1189 },
1190 old_from_new
1191 );
1192 }
1193
1194 #[test]
1195 fn test_region_change_backward_compatibility() {
1196 let region_change_json = r#"{
1198 "metadata": {
1199 "column_metadatas": [
1200 {"column_schema":{"name":"a","data_type":{"Int64":{}},"is_nullable":false,"is_time_index":false,"default_constraint":null,"metadata":{}},"semantic_type":"Tag","column_id":1},
1201 {"column_schema":{"name":"b","data_type":{"Int64":{}},"is_nullable":false,"is_time_index":false,"default_constraint":null,"metadata":{}},"semantic_type":"Field","column_id":2},
1202 {"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}
1203 ],
1204 "primary_key": [
1205 1
1206 ],
1207 "region_id": 42,
1208 "schema_version": 0
1209 }
1210 }"#;
1211
1212 let region_change: RegionChange = serde_json::from_str(region_change_json).unwrap();
1213 assert_eq!(region_change.sst_format, FormatType::PrimaryKey);
1214
1215 let region_change = RegionChange {
1217 metadata: region_change.metadata.clone(),
1218 sst_format: FormatType::Flat,
1219 append_mode: None,
1220 };
1221
1222 let serialized = serde_json::to_string(®ion_change).unwrap();
1223 let deserialized: RegionChange = serde_json::from_str(&serialized).unwrap();
1224 assert_eq!(deserialized.sst_format, FormatType::Flat);
1225 }
1226
1227 #[test]
1228 fn test_removed_file_compatibility() {
1229 let file_id = FileId::random();
1230 let json_str = format!("\"{}\"", file_id);
1232 let removed_file: RemovedFile = serde_json::from_str(&json_str).unwrap();
1233 assert_eq!(removed_file, RemovedFile::File(file_id, None));
1234
1235 let removed_file_v2 = RemovedFile::File(file_id, Some(10));
1237 let json_v2 = serde_json::to_string(&removed_file_v2).unwrap();
1238 let deserialized_v2: RemovedFile = serde_json::from_str(&json_v2).unwrap();
1239 assert_eq!(removed_file_v2, deserialized_v2);
1240
1241 let removed_index = RemovedFile::Index(file_id, 20);
1243 let json_index = serde_json::to_string(&removed_index).unwrap();
1244 let deserialized_index: RemovedFile = serde_json::from_str(&json_index).unwrap();
1245 assert_eq!(removed_index, deserialized_index);
1246
1247 let removed_file = RemovedFile::File(file_id, None);
1249 let json = serde_json::to_string(&removed_file).unwrap();
1250 let deserialized: RemovedFile = serde_json::from_str(&json).unwrap();
1251 assert_eq!(removed_file, deserialized);
1252
1253 let json_set = format!("[\"{}\"]", file_id);
1259 let removed_files_set: HashSet<RemovedFile> = serde_json::from_str(&json_set).unwrap();
1260 assert!(removed_files_set.contains(&RemovedFile::File(file_id, None)));
1261 }
1262
1263 #[test]
1286 fn test_removed_files_backward_compatibility() {
1287 #[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
1289 struct OldRemovedFiles {
1290 pub removed_at: i64,
1291 pub file_ids: HashSet<FileId>,
1292 }
1293
1294 let mut file_ids = HashSet::new();
1296 file_ids.insert(FileId::random());
1297 file_ids.insert(FileId::random());
1298
1299 let old_removed_files = OldRemovedFiles {
1300 removed_at: 1234567890,
1301 file_ids,
1302 };
1303
1304 let old_json = serde_json::to_string(&old_removed_files).unwrap();
1306
1307 let result: Result<RemovedFiles, _> = serde_json::from_str(&old_json);
1309
1310 assert!(result.is_ok(), "{:?}", result);
1312 let removed_files = result.unwrap();
1313 assert_eq!(removed_files.removed_at, 1234567890);
1314 assert!(removed_files.files.is_empty());
1315
1316 let file_id = FileId::random();
1318 let new_json = format!(
1319 r#"{{
1320 "removed_at": 1234567890,
1321 "files": ["{}"]
1322 }}"#,
1323 file_id
1324 );
1325
1326 let result: Result<RemovedFiles, _> = serde_json::from_str(&new_json);
1327 assert!(result.is_ok());
1328 let removed_files = result.unwrap();
1329 assert_eq!(removed_files.removed_at, 1234567890);
1330 assert_eq!(removed_files.files.len(), 1);
1331 assert!(
1332 removed_files
1333 .files
1334 .contains(&RemovedFile::File(file_id, None))
1335 );
1336 }
1337}