Skip to main content

mito2/region/
version.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Version control of mito engine.
16//!
17//! Version is an immutable snapshot of region's metadata.
18//!
19//! To read latest data from `VersionControl`, we should
20//! 1. Acquire `Version` from `VersionControl`.
21//! 2. Then acquire last sequence.
22//!
23//! Reason: data may be flushed/compacted and some data with old sequence may be removed
24//! and became invisible between step 1 and 2, so need to acquire version at first.
25
26use std::sync::{Arc, RwLock};
27use std::time::Duration;
28
29use common_telemetry::info;
30use store_api::metadata::RegionMetadataRef;
31use store_api::storage::{RegionId, SequenceNumber};
32
33use crate::error::Result;
34use crate::manifest::action::{RegionEdit, TruncateKind};
35use crate::memtable::time_partition::{TimePartitions, TimePartitionsRef};
36use crate::memtable::version::{MemtableVersion, MemtableVersionRef};
37use crate::memtable::{MemtableBuilderRef, MemtableId};
38use crate::region::options::RegionOptions;
39use crate::sst::file::FileMeta;
40use crate::sst::file_purger::FilePurgerRef;
41use crate::sst::version::{SstVersion, SstVersionRef};
42use crate::wal::EntryId;
43
44/// Controls metadata and sequence numbers for a region.
45///
46/// It manages metadata in a copy-on-write fashion. Any modification to a region's metadata
47/// will generate a new [Version].
48#[derive(Debug)]
49pub(crate) struct VersionControl {
50    data: RwLock<VersionControlData>,
51}
52
53impl VersionControl {
54    /// Returns a new [VersionControl] with specific `version`.
55    pub(crate) fn new(version: Version) -> VersionControl {
56        // Initialize sequence and entry id from flushed sequence and entry id.
57        let (flushed_sequence, flushed_entry_id) =
58            (version.flushed_sequence, version.flushed_entry_id);
59        VersionControl {
60            data: RwLock::new(VersionControlData {
61                version: Arc::new(version),
62                committed_sequence: flushed_sequence,
63                last_entry_id: flushed_entry_id,
64                is_dropped: false,
65            }),
66        }
67    }
68
69    /// Returns the id of the region controlled by this instance.
70    pub(crate) fn region_id(&self) -> RegionId {
71        self.data.read().unwrap().version.metadata.region_id
72    }
73
74    /// Returns current copy of data.
75    pub(crate) fn current(&self) -> VersionControlData {
76        self.data.read().unwrap().clone()
77    }
78
79    /// Updates the `committed_sequence` of version.
80    pub(crate) fn set_committed_sequence(&self, seq: SequenceNumber) {
81        let mut data = self.data.write().unwrap();
82        data.committed_sequence = seq;
83    }
84
85    /// Updates committed sequence and entry id.
86    pub(crate) fn set_sequence_and_entry_id(&self, seq: SequenceNumber, entry_id: EntryId) {
87        let mut data = self.data.write().unwrap();
88        data.committed_sequence = seq;
89        data.last_entry_id = entry_id;
90    }
91
92    /// Updates last entry id.
93    pub(crate) fn set_entry_id(&self, entry_id: EntryId) {
94        let mut data = self.data.write().unwrap();
95        data.last_entry_id = entry_id;
96    }
97
98    /// Sequence number of last committed data.
99    pub(crate) fn committed_sequence(&self) -> SequenceNumber {
100        self.data.read().unwrap().committed_sequence
101    }
102
103    /// Freezes the mutable memtable if it is not empty.
104    pub(crate) fn freeze_mutable(&self) -> Result<()> {
105        let version = self.current().version;
106        let time_window = version.compaction_time_window;
107
108        let Some(new_memtables) = version
109            .memtables
110            .freeze_mutable(&version.metadata, time_window)?
111        else {
112            return Ok(());
113        };
114
115        // Create a new version with memtable switched.
116        let new_version = Arc::new(
117            VersionBuilder::from_version(version)
118                .memtables(new_memtables)
119                .build(),
120        );
121
122        let mut version_data = self.data.write().unwrap();
123        version_data.version = new_version;
124
125        Ok(())
126    }
127
128    /// Applies region option changes and generates a new version.
129    pub(crate) fn alter_options(&self, options: RegionOptions) {
130        let version = self.current().version;
131        let new_version = Arc::new(
132            VersionBuilder::from_version(version)
133                .options(options)
134                .build(),
135        );
136        let mut version_data = self.data.write().unwrap();
137        version_data.version = new_version;
138    }
139
140    /// Apply edit to current version.
141    ///
142    /// If `edit` is None, only removes the specified memtables.
143    pub(crate) fn apply_edit(
144        &self,
145        edit: Option<RegionEdit>,
146        memtables_to_remove: &[MemtableId],
147        purger: FilePurgerRef,
148    ) {
149        let version = self.current().version;
150        let builder = VersionBuilder::from_version(version);
151        let committed_sequence = edit.as_ref().and_then(|e| e.committed_sequence);
152        let builder = if let Some(edit) = edit {
153            builder.apply_edit(edit, purger)
154        } else {
155            builder
156        };
157        let new_version = Arc::new(builder.remove_memtables(memtables_to_remove).build());
158
159        let mut version_data = self.data.write().unwrap();
160        version_data.committed_sequence = if let Some(committed_in_edit) = committed_sequence {
161            version_data.committed_sequence.max(committed_in_edit)
162        } else {
163            version_data.committed_sequence
164        };
165        version_data.version = new_version;
166    }
167
168    /// Mark all opened files as deleted and set the delete marker in [VersionControlData]
169    pub(crate) fn mark_dropped(&self) {
170        let version = self.current().version;
171        let memtable_builder = version.memtables.mutable.memtable_builder().clone();
172        let new_mutable =
173            Self::new_mutable_from_version(&version, version.metadata.clone(), memtable_builder);
174
175        let mut data = self.data.write().unwrap();
176        data.is_dropped = true;
177        data.version.ssts.mark_all_deleted();
178        // Reset version so we can release the reference to memtables and SSTs.
179        let new_version =
180            Arc::new(VersionBuilder::new(version.metadata.clone(), new_mutable).build());
181        data.version = new_version;
182    }
183
184    /// Alter schema of the region.
185    ///
186    /// It replaces existing mutable memtable with a memtable that uses the
187    /// new schema. Memtables of the version must be empty.
188    pub(crate) fn alter_schema(&self, metadata: RegionMetadataRef) {
189        let version = self.current().version;
190        let memtable_builder = version.memtables.mutable.memtable_builder().clone();
191        let new_mutable =
192            Self::new_mutable_from_version(&version, metadata.clone(), memtable_builder);
193        debug_assert!(version.memtables.mutable.is_empty());
194        debug_assert!(version.memtables.immutables().is_empty());
195        let new_version = Arc::new(
196            VersionBuilder::from_version(version)
197                .metadata(metadata)
198                .memtables(MemtableVersion::new(new_mutable))
199                .build(),
200        );
201
202        let mut version_data = self.data.write().unwrap();
203        version_data.version = new_version;
204    }
205
206    /// Alter metadata of the region without rebuilding memtables.
207    pub(crate) fn alter_metadata(&self, metadata: RegionMetadataRef) {
208        let version = self.current().version;
209        let new_version = Arc::new(
210            VersionBuilder::from_version(version)
211                .metadata(metadata)
212                .build(),
213        );
214
215        let mut version_data = self.data.write().unwrap();
216        version_data.version = new_version;
217    }
218
219    /// Alter schema and format of the region.
220    ///
221    /// It replaces existing mutable memtable with a memtable that uses the
222    /// new format. Memtables of the version must be empty.
223    pub(crate) fn alter_schema_and_format(
224        &self,
225        metadata: RegionMetadataRef,
226        options: RegionOptions,
227        memtable_builder: MemtableBuilderRef,
228    ) {
229        let version = self.current().version;
230        // Use the new metadata to build `TimePartitions`.
231        let new_mutable =
232            Self::new_mutable_from_version(&version, metadata.clone(), memtable_builder);
233        debug_assert!(version.memtables.mutable.is_empty());
234        debug_assert!(version.memtables.immutables().is_empty());
235        let new_version = Arc::new(
236            VersionBuilder::from_version(version)
237                .metadata(metadata)
238                .options(options)
239                .memtables(MemtableVersion::new(new_mutable))
240                .build(),
241        );
242
243        let mut version_data = self.data.write().unwrap();
244        version_data.version = new_version;
245    }
246
247    /// Truncate current version.
248    pub(crate) fn truncate(&self, truncate_kind: TruncateKind) {
249        let version = self.current().version;
250
251        match truncate_kind {
252            TruncateKind::All {
253                truncated_entry_id,
254                truncated_sequence,
255            } => {
256                let memtable_builder = version.memtables.mutable.memtable_builder().clone();
257                let new_mutable = Self::new_mutable_from_version(
258                    &version,
259                    version.metadata.clone(),
260                    memtable_builder,
261                );
262                let new_version = Arc::new(
263                    VersionBuilder::from_version(version)
264                        .memtables(MemtableVersion::new(new_mutable))
265                        .clear_files()
266                        .flushed_entry_id(truncated_entry_id)
267                        .flushed_sequence(truncated_sequence)
268                        .truncated_entry_id(Some(truncated_entry_id))
269                        .build(),
270                );
271
272                let mut version_data = self.data.write().unwrap();
273                version_data.version.ssts.mark_all_deleted();
274                version_data.version = new_version;
275            }
276            TruncateKind::Partial { files_to_remove } => {
277                let new_version = Arc::new(
278                    VersionBuilder::from_version(version)
279                        .remove_files(files_to_remove.into_iter())
280                        .build(),
281                );
282
283                let mut version_data = self.data.write().unwrap();
284                // notice since it's partial, no need to mark all files as deleted
285                version_data.version = new_version;
286            }
287        };
288    }
289
290    /// Discards all memtables while preserving persisted SST files.
291    pub(crate) fn discard_unflushed(
292        &self,
293        discarded_entry_id: EntryId,
294        discarded_sequence: SequenceNumber,
295    ) {
296        let version = self.current().version;
297        let memtable_builder = version.memtables.mutable.memtable_builder().clone();
298        let new_mutable =
299            Self::new_mutable_from_version(&version, version.metadata.clone(), memtable_builder);
300        let new_version = Arc::new(
301            VersionBuilder::from_version(version)
302                .memtables(MemtableVersion::new(new_mutable))
303                .flushed_entry_id(discarded_entry_id)
304                .flushed_sequence(discarded_sequence)
305                .build(),
306        );
307
308        let mut version_data = self.data.write().unwrap();
309        version_data.version = new_version;
310    }
311
312    /// Overwrites the current version with a new version.
313    pub(crate) fn overwrite_current(&self, version: VersionRef) {
314        let mut version_data = self.data.write().unwrap();
315        version_data.version = version;
316    }
317
318    fn new_mutable_from_version(
319        version: &Version,
320        metadata: RegionMetadataRef,
321        memtable_builder: MemtableBuilderRef,
322    ) -> TimePartitionsRef {
323        Arc::new(TimePartitions::new(
324            metadata,
325            memtable_builder,
326            version.memtables.mutable.next_memtable_id(),
327            Some(version.memtables.mutable.part_duration()),
328        ))
329    }
330}
331
332pub(crate) type VersionControlRef = Arc<VersionControl>;
333
334/// Data of [VersionControl].
335#[derive(Debug, Clone)]
336pub(crate) struct VersionControlData {
337    /// Latest version.
338    pub(crate) version: VersionRef,
339    /// Sequence number of last committed data.
340    ///
341    /// Starts from 1 (zero means no data).
342    pub(crate) committed_sequence: SequenceNumber,
343    /// Last WAL entry Id.
344    ///
345    /// Starts from 1 (zero means no data).
346    pub(crate) last_entry_id: EntryId,
347    /// Marker of whether this region is dropped/dropping
348    pub(crate) is_dropped: bool,
349}
350
351impl VersionControlData {
352    /// Approximate timeseries count in current version.
353    pub(crate) fn series_count(&self) -> usize {
354        self.version.memtables.mutable.series_count()
355    }
356}
357
358/// Static metadata of a region.
359#[derive(Clone, Debug)]
360pub(crate) struct Version {
361    /// Metadata of the region.
362    ///
363    /// Altering metadata isn't frequent, storing metadata in Arc to allow sharing
364    /// metadata and reuse metadata when creating a new `Version`.
365    pub(crate) metadata: RegionMetadataRef,
366    /// Mutable and immutable memtables.
367    ///
368    /// Wrapped in Arc to make clone of `Version` much cheaper.
369    pub(crate) memtables: MemtableVersionRef,
370    /// SSTs of the region.
371    pub(crate) ssts: SstVersionRef,
372    /// Inclusive max WAL entry id of flushed data.
373    pub(crate) flushed_entry_id: EntryId,
374    /// Inclusive max sequence of flushed data.
375    pub(crate) flushed_sequence: SequenceNumber,
376    /// Latest entry id during the truncating table.
377    ///
378    /// Used to check if it is a flush task during the truncating table.
379    pub(crate) truncated_entry_id: Option<EntryId>,
380    /// Inferred compaction time window from flush.
381    ///
382    /// If compaction options contain a time window, it will overwrite this value
383    /// when creating a new version from the [VersionBuilder].
384    pub(crate) compaction_time_window: Option<Duration>,
385    /// Options of the region.
386    pub(crate) options: RegionOptions,
387}
388
389pub(crate) type VersionRef = Arc<Version>;
390
391/// Version builder.
392pub(crate) struct VersionBuilder {
393    metadata: RegionMetadataRef,
394    memtables: MemtableVersionRef,
395    ssts: SstVersionRef,
396    flushed_entry_id: EntryId,
397    flushed_sequence: SequenceNumber,
398    truncated_entry_id: Option<EntryId>,
399    compaction_time_window: Option<Duration>,
400    options: RegionOptions,
401}
402
403impl VersionBuilder {
404    /// Returns a new builder.
405    pub(crate) fn new(metadata: RegionMetadataRef, mutable: TimePartitionsRef) -> Self {
406        VersionBuilder {
407            metadata,
408            memtables: Arc::new(MemtableVersion::new(mutable)),
409            ssts: Arc::new(SstVersion::new()),
410            flushed_entry_id: 0,
411            flushed_sequence: 0,
412            truncated_entry_id: None,
413            compaction_time_window: None,
414            options: RegionOptions::default(),
415        }
416    }
417
418    /// Returns a new builder from an existing version.
419    pub(crate) fn from_version(version: VersionRef) -> Self {
420        VersionBuilder {
421            metadata: version.metadata.clone(),
422            memtables: version.memtables.clone(),
423            ssts: version.ssts.clone(),
424            flushed_entry_id: version.flushed_entry_id,
425            flushed_sequence: version.flushed_sequence,
426            truncated_entry_id: version.truncated_entry_id,
427            compaction_time_window: version.compaction_time_window,
428            options: version.options.clone(),
429        }
430    }
431
432    /// Sets memtables.
433    pub(crate) fn memtables(mut self, memtables: MemtableVersion) -> Self {
434        self.memtables = Arc::new(memtables);
435        self
436    }
437
438    /// Sets metadata.
439    pub(crate) fn metadata(mut self, metadata: RegionMetadataRef) -> Self {
440        self.metadata = metadata;
441        self
442    }
443
444    /// Sets flushed entry id.
445    pub(crate) fn flushed_entry_id(mut self, entry_id: EntryId) -> Self {
446        self.flushed_entry_id = entry_id;
447        self
448    }
449
450    /// Sets flushed sequence.
451    pub(crate) fn flushed_sequence(mut self, sequence: SequenceNumber) -> Self {
452        self.flushed_sequence = sequence;
453        self
454    }
455
456    /// Sets truncated entry id.
457    pub(crate) fn truncated_entry_id(mut self, entry_id: Option<EntryId>) -> Self {
458        self.truncated_entry_id = entry_id;
459        self
460    }
461
462    /// Sets compaction time window.
463    pub(crate) fn compaction_time_window(mut self, window: Option<Duration>) -> Self {
464        self.compaction_time_window = window;
465        self
466    }
467
468    /// Sets options.
469    pub(crate) fn options(mut self, options: RegionOptions) -> Self {
470        self.options = options;
471        self
472    }
473
474    /// Apply edit to the builder.
475    pub(crate) fn apply_edit(mut self, edit: RegionEdit, file_purger: FilePurgerRef) -> Self {
476        if let Some(entry_id) = edit.flushed_entry_id {
477            self.flushed_entry_id = self.flushed_entry_id.max(entry_id);
478        }
479        if let Some(sequence) = edit.flushed_sequence {
480            self.flushed_sequence = self.flushed_sequence.max(sequence);
481        }
482        if let Some(window) = edit.compaction_time_window {
483            self.compaction_time_window = Some(window);
484        }
485        if !edit.files_to_add.is_empty() || !edit.files_to_remove.is_empty() {
486            let mut ssts = (*self.ssts).clone();
487            ssts.add_files(file_purger, edit.files_to_add.into_iter());
488            ssts.remove_files(edit.files_to_remove.into_iter());
489            self.ssts = Arc::new(ssts);
490        }
491
492        self
493    }
494
495    /// Remove memtables from the builder.
496    pub(crate) fn remove_memtables(mut self, ids: &[MemtableId]) -> Self {
497        if !ids.is_empty() {
498            let mut memtables = (*self.memtables).clone();
499            memtables.remove_memtables(ids);
500            self.memtables = Arc::new(memtables);
501        }
502        self
503    }
504
505    /// Add files to the builder.
506    pub(crate) fn add_files(
507        mut self,
508        file_purger: FilePurgerRef,
509        files: impl Iterator<Item = FileMeta>,
510    ) -> Self {
511        let mut ssts = (*self.ssts).clone();
512        ssts.add_files(file_purger, files);
513        self.ssts = Arc::new(ssts);
514
515        self
516    }
517
518    pub(crate) fn remove_files(mut self, files: impl Iterator<Item = FileMeta>) -> Self {
519        let mut ssts = (*self.ssts).clone();
520        ssts.remove_files(files);
521        self.ssts = Arc::new(ssts);
522
523        self
524    }
525
526    /// Clear all files in the builder.
527    pub(crate) fn clear_files(mut self) -> Self {
528        self.ssts = Arc::new(SstVersion::new());
529        self
530    }
531
532    /// Builds a new [Version] from the builder.
533    /// It overwrites the window size by compaction option.
534    pub(crate) fn build(self) -> Version {
535        let compaction_time_window = self
536            .options
537            .compaction
538            .time_window()
539            .or(self.compaction_time_window);
540        if self.compaction_time_window.is_some()
541            && compaction_time_window != self.compaction_time_window
542        {
543            info!(
544                "VersionBuilder overwrites region compaction time window from {:?} to {:?}, region: {}",
545                self.compaction_time_window, compaction_time_window, self.metadata.region_id
546            );
547        }
548
549        Version {
550            metadata: self.metadata,
551            memtables: self.memtables,
552            ssts: self.ssts,
553            flushed_entry_id: self.flushed_entry_id,
554            flushed_sequence: self.flushed_sequence,
555            truncated_entry_id: self.truncated_entry_id,
556            compaction_time_window,
557            options: self.options,
558        }
559    }
560}