1use 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#[derive(Debug)]
49pub(crate) struct VersionControl {
50 data: RwLock<VersionControlData>,
51}
52
53impl VersionControl {
54 pub(crate) fn new(version: Version) -> VersionControl {
56 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 pub(crate) fn region_id(&self) -> RegionId {
71 self.data.read().unwrap().version.metadata.region_id
72 }
73
74 pub(crate) fn current(&self) -> VersionControlData {
76 self.data.read().unwrap().clone()
77 }
78
79 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 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 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 pub(crate) fn committed_sequence(&self) -> SequenceNumber {
100 self.data.read().unwrap().committed_sequence
101 }
102
103 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 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 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 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 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 let new_version =
180 Arc::new(VersionBuilder::new(version.metadata.clone(), new_mutable).build());
181 data.version = new_version;
182 }
183
184 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 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 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 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 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 version_data.version = new_version;
286 }
287 };
288 }
289
290 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 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#[derive(Debug, Clone)]
336pub(crate) struct VersionControlData {
337 pub(crate) version: VersionRef,
339 pub(crate) committed_sequence: SequenceNumber,
343 pub(crate) last_entry_id: EntryId,
347 pub(crate) is_dropped: bool,
349}
350
351impl VersionControlData {
352 pub(crate) fn series_count(&self) -> usize {
354 self.version.memtables.mutable.series_count()
355 }
356}
357
358#[derive(Clone, Debug)]
360pub(crate) struct Version {
361 pub(crate) metadata: RegionMetadataRef,
366 pub(crate) memtables: MemtableVersionRef,
370 pub(crate) ssts: SstVersionRef,
372 pub(crate) flushed_entry_id: EntryId,
374 pub(crate) flushed_sequence: SequenceNumber,
376 pub(crate) truncated_entry_id: Option<EntryId>,
380 pub(crate) compaction_time_window: Option<Duration>,
385 pub(crate) options: RegionOptions,
387}
388
389pub(crate) type VersionRef = Arc<Version>;
390
391pub(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 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 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 pub(crate) fn memtables(mut self, memtables: MemtableVersion) -> Self {
434 self.memtables = Arc::new(memtables);
435 self
436 }
437
438 pub(crate) fn metadata(mut self, metadata: RegionMetadataRef) -> Self {
440 self.metadata = metadata;
441 self
442 }
443
444 pub(crate) fn flushed_entry_id(mut self, entry_id: EntryId) -> Self {
446 self.flushed_entry_id = entry_id;
447 self
448 }
449
450 pub(crate) fn flushed_sequence(mut self, sequence: SequenceNumber) -> Self {
452 self.flushed_sequence = sequence;
453 self
454 }
455
456 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 pub(crate) fn compaction_time_window(mut self, window: Option<Duration>) -> Self {
464 self.compaction_time_window = window;
465 self
466 }
467
468 pub(crate) fn options(mut self, options: RegionOptions) -> Self {
470 self.options = options;
471 self
472 }
473
474 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 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 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 pub(crate) fn clear_files(mut self) -> Self {
528 self.ssts = Arc::new(SstVersion::new());
529 self
530 }
531
532 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}