Skip to main content

mito2/worker/
handle_manifest.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//! Handles manifest.
16//!
17//! It updates the manifest and applies the changes to the region in background.
18
19use std::collections::{HashMap, VecDeque};
20use std::num::NonZeroU64;
21use std::sync::Arc;
22
23use common_telemetry::{debug, info, warn};
24use parquet::file::metadata::PageIndexPolicy;
25use snafu::ResultExt;
26use store_api::logstore::LogStore;
27use store_api::metadata::RegionMetadataRef;
28use store_api::storage::{RegionId, SequenceNumber};
29
30use crate::cache::CacheManagerRef;
31use crate::cache::file_cache::{FileType, IndexKey};
32use crate::config::IndexBuildMode;
33use crate::error::{EditRegionSnafu, RegionBusySnafu, RegionNotFoundSnafu, Result};
34use crate::manifest::action::{
35    RegionChange, RegionEdit, RegionMetaAction, RegionMetaActionList, RegionTruncate,
36};
37use crate::memtable::MemtableBuilderProvider;
38use crate::metrics::WRITE_CACHE_INFLIGHT_DOWNLOAD;
39use crate::region::opener::{sanitize_region_options, version_builder_from_manifest};
40use crate::region::options::RegionOptions;
41use crate::region::version::VersionControlRef;
42use crate::region::{MitoRegionRef, RegionLeaderState, RegionRoleState};
43use crate::request::{
44    BackgroundNotify, BuildIndexRequest, DiscardUnflushedResult, OptionOutputTx,
45    RegionChangeResult, RegionEditRequest, RegionEditResult, RegionSyncRequest, TruncateResult,
46    WorkerRequest, WorkerRequestWithTime,
47};
48use crate::sst::index::IndexBuildType;
49use crate::sst::location;
50use crate::wal::EntryId;
51use crate::worker::{RegionWorkerLoop, WorkerListener};
52
53pub(crate) type RegionEditQueues = HashMap<RegionId, RegionEditQueue>;
54
55/// A queue for region edit requests received while the region is already `Editing`.
56///
57/// Normal writes and bulk inserts that arrive during `Editing` use the stalled-write queue instead.
58/// When an edit completes, those writes are handled before the next queued edit starts, preserving
59/// sequence ordering between direct SST edits and WAL/memtable writes unless global reject
60/// backpressure rejects them first.
61/// Everything is done in the region worker loop.
62pub(crate) struct RegionEditQueue {
63    region_id: RegionId,
64    requests: VecDeque<RegionEditRequest>,
65}
66
67impl RegionEditQueue {
68    const QUEUE_MAX_LEN: usize = 128;
69
70    fn new(region_id: RegionId) -> Self {
71        Self {
72            region_id,
73            requests: VecDeque::new(),
74        }
75    }
76
77    fn enqueue(&mut self, request: RegionEditRequest) {
78        if self.requests.len() > Self::QUEUE_MAX_LEN {
79            request.waiters.reply_with(|| {
80                RegionBusySnafu {
81                    region_id: self.region_id,
82                }
83                .fail()
84            });
85            return;
86        };
87        self.requests.push_back(request);
88    }
89
90    fn dequeue(&mut self) -> Option<RegionEditRequest> {
91        fn can_merge(edit: &RegionEdit) -> bool {
92            // Only the `RegionEdit`:
93            // 1. contains the "raw" (file without a sequence) files to add,
94            // 2. and no `committed_sequence`,
95            // 3. and all other fields are empty,
96            // can it be merged.
97            //
98            // However, merging them means they will all share a same sequence, and if there are
99            // overlapping data in the files, the dedup is uncertain. This is a caution that must
100            // be noticed for editing region.
101            edit.files_to_add.iter().all(|f| f.sequence.is_none())
102                && edit.files_to_remove.is_empty()
103                && edit.timestamp_ms.is_none()
104                && edit.compaction_time_window.is_none()
105                && edit.flushed_entry_id.is_none()
106                && edit.flushed_sequence.is_none()
107                && edit.committed_sequence.is_none()
108        }
109
110        let mut merged = self.requests.pop_front()?;
111        if !can_merge(&merged.edit) {
112            return Some(merged);
113        }
114
115        while let Some(request) = self
116            .requests
117            .pop_front_if(|request| can_merge(&request.edit))
118        {
119            merged.edit.files_to_add.extend(request.edit.files_to_add);
120            merged.waiters.merge(request.waiters);
121        }
122        debug!(
123            "the files to add: [{}] are merged in one edit",
124            merged
125                .edit
126                .files_to_add
127                .iter()
128                .map(|x| x.file_id.to_string())
129                .collect::<Vec<_>>()
130                .join(", ")
131        );
132        Some(merged)
133    }
134
135    fn is_empty(&self) -> bool {
136        self.requests.is_empty()
137    }
138
139    fn reject_all_as_not_found(mut self) {
140        while let Some(request) = self.requests.pop_front() {
141            request.waiters.reply_with(|| {
142                RegionNotFoundSnafu {
143                    region_id: self.region_id,
144                }
145                .fail()
146            });
147        }
148    }
149}
150
151impl<S: LogStore> RegionWorkerLoop<S> {
152    /// Rejects queued region edit requests as region not found.
153    pub(crate) fn reject_region_edit_queue_as_not_found(&mut self, region_id: RegionId) {
154        if let Some(edit_queue) = self.region_edit_queues.remove(&region_id) {
155            edit_queue.reject_all_as_not_found();
156        }
157    }
158
159    /// Handles region change result.
160    pub(crate) async fn handle_manifest_region_change_result(
161        &mut self,
162        change_result: RegionChangeResult,
163    ) {
164        let region = match self.regions.get_region(change_result.region_id) {
165            Some(region) => region,
166            None => {
167                self.reject_region_stalled_requests(&change_result.region_id);
168                change_result.sender.send(
169                    RegionNotFoundSnafu {
170                        region_id: change_result.region_id,
171                    }
172                    .fail(),
173                );
174                return;
175            }
176        };
177
178        if change_result.result.is_ok() {
179            // Updates the region metadata and format.
180            Self::update_region_version(
181                &region.version_control,
182                change_result.new_meta,
183                change_result.new_options,
184                &self.memtable_builder_provider,
185            );
186        }
187
188        // Sets the region as writable.
189        region.switch_state_to_writable(RegionLeaderState::Altering);
190        // Sends the result.
191        change_result.sender.send(change_result.result.map(|_| 0));
192
193        // In async mode, rebuild index after index metadata changed.
194        if self.config.index.build_mode == IndexBuildMode::Async && change_result.need_index {
195            self.handle_rebuild_index(
196                BuildIndexRequest {
197                    region_id: region.region_id,
198                    build_type: IndexBuildType::SchemaChange,
199                    file_metas: Vec::new(),
200                },
201                OptionOutputTx::new(None),
202            )
203            .await;
204        }
205        // Handles the stalled requests.
206        self.handle_region_stalled_requests(&change_result.region_id, true)
207            .await;
208    }
209
210    /// Handles region sync request.
211    ///
212    /// Updates the manifest to at least the given version.
213    /// **Note**: The installed version may be greater than the given version.
214    pub(crate) async fn handle_region_sync(&mut self, request: RegionSyncRequest) {
215        let region_id = request.region_id;
216        let sender = request.sender;
217        let region = match self.regions.follower_region(region_id) {
218            Ok(region) => region,
219            Err(e) => {
220                let _ = sender.send(Err(e));
221                return;
222            }
223        };
224
225        let original_manifest_version = region.manifest_ctx.manifest_version().await;
226        let manifest = match region
227            .manifest_ctx
228            .install_manifest_to(request.manifest_version)
229            .await
230        {
231            Ok(manifest) => manifest,
232            Err(e) => {
233                let _ = sender.send(Err(e));
234                return;
235            }
236        };
237        let version = region.version();
238        let mut region_options = version.options.clone();
239        let old_format = region_options.sst_format.unwrap_or_default();
240        // Updates the region options with the manifest.
241        sanitize_region_options(&manifest, &mut region_options);
242        if !version.memtables.is_empty() {
243            let current = region.version_control.current();
244            warn!(
245                "Region {} memtables is not empty, which should not happen, manifest version: {}, last entry id: {}",
246                region.region_id, manifest.manifest_version, current.last_entry_id
247            );
248        }
249
250        // We should sanitize the region options before creating a new memtable.
251        let memtable_builder = if old_format != region_options.sst_format.unwrap_or_default() {
252            // Format changed, also needs to replace the memtable builder.
253            Some(
254                self.memtable_builder_provider
255                    .builder_for_options(&region_options),
256            )
257        } else {
258            None
259        };
260        let new_mutable = Arc::new(
261            region
262                .version()
263                .memtables
264                .mutable
265                .new_with_part_duration(version.compaction_time_window, memtable_builder),
266        );
267        // Here it assumes the leader has backfilled the partition_expr of the metadata.
268        let metadata = manifest.metadata.clone();
269
270        let version_builder = version_builder_from_manifest(
271            &manifest,
272            metadata,
273            region.file_purger.clone(),
274            new_mutable,
275            region_options,
276        );
277        let version = version_builder.build();
278        region.version_control.overwrite_current(Arc::new(version));
279
280        let updated = manifest.manifest_version > original_manifest_version;
281        let _ = sender.send(Ok((manifest.manifest_version, updated)));
282    }
283}
284
285impl<S: LogStore> RegionWorkerLoop<S> {
286    /// Handles region edit request.
287    pub(crate) fn handle_region_edit(&mut self, request: RegionEditRequest) {
288        let region_id = request.region_id;
289        let Some(region) = self.regions.get_region(region_id) else {
290            request
291                .waiters
292                .reply_with(|| RegionNotFoundSnafu { region_id }.fail());
293            return;
294        };
295
296        if !region.is_writable() {
297            if region.state() == RegionRoleState::Leader(RegionLeaderState::Editing) {
298                self.region_edit_queues
299                    .entry(region_id)
300                    .or_insert_with(|| RegionEditQueue::new(region_id))
301                    .enqueue(request);
302            } else {
303                request
304                    .waiters
305                    .reply_with(|| RegionBusySnafu { region_id }.fail());
306            }
307            return;
308        }
309
310        let RegionEditRequest {
311            region_id: _,
312            mut edit,
313            waiters,
314            preload_sst_cache,
315        } = request;
316        let file_sequence = region.version_control.committed_sequence() + 1;
317        edit.committed_sequence = Some(file_sequence);
318
319        // Generic region edits (direct/import/staging) assign a new destination
320        // sequence domain but cannot prove or rewrite the physical per-row sequence
321        // column. The added files are therefore atomic/untrusted: clear the
322        // `preserve_row_sequence` marker so exact sequence-range scans fail closed
323        // until the scan reaches the assigned max sequence.
324        for file in &mut edit.files_to_add {
325            file.sequence = NonZeroU64::new(file_sequence);
326            file.preserve_row_sequence = false;
327        }
328
329        // Allow retrieving `is_staging` before spawn the edit region task.
330        let is_staging = region.is_staging();
331        let expect_state = if is_staging {
332            RegionLeaderState::Staging
333        } else {
334            RegionLeaderState::Writable
335        };
336        // Marks the region as editing.
337        if let Err(e) = region.set_editing(expect_state) {
338            let e = Arc::new(e);
339            waiters.reply_with(|| Err(e.clone()).context(EditRegionSnafu { region_id }));
340            return;
341        }
342
343        let request_sender = self.sender.clone();
344        let cache_manager = self.cache_manager.clone();
345        let listener = self.listener.clone();
346        // Now the region is in editing state.
347        // Updates manifest in background.
348        common_runtime::spawn_global(async move {
349            let result = edit_region(
350                &region,
351                edit.clone(),
352                cache_manager,
353                listener,
354                is_staging,
355                preload_sst_cache,
356            )
357            .await
358            .map_err(Arc::new);
359            let notify = WorkerRequest::Background {
360                region_id,
361                notify: BackgroundNotify::RegionEdit(RegionEditResult {
362                    region_id,
363                    waiters,
364                    edit,
365                    result,
366                    // we always need to restore region state after region edit
367                    update_region_state: true,
368                    is_staging,
369                }),
370            };
371
372            // We don't set state back as the worker loop is already exited.
373            if let Err(res) = request_sender
374                .send(WorkerRequestWithTime::new(notify))
375                .await
376            {
377                warn!(
378                    "Failed to send region edit result back to the worker, region_id: {}, res: {:?}",
379                    region_id, res
380                );
381            }
382        });
383    }
384
385    /// Handles region edit result.
386    pub(crate) async fn handle_region_edit_result(&mut self, edit_result: RegionEditResult) {
387        let region = match self.regions.get_region(edit_result.region_id) {
388            Some(region) => region,
389            None => {
390                // Fail writes stalled behind this edit if the region was removed before the
391                // edit-completion notification reached the worker.
392                self.fail_region_stalled_requests_as_not_found(&edit_result.region_id);
393                self.reject_region_edit_queue_as_not_found(edit_result.region_id);
394
395                edit_result.waiters.reply_with(|| {
396                    RegionNotFoundSnafu {
397                        region_id: edit_result.region_id,
398                    }
399                    .fail()
400                });
401                return;
402            }
403        };
404
405        let need_compaction = if edit_result.is_staging {
406            if edit_result.update_region_state {
407                // For staging regions, edits are not applied immediately,
408                // as they remain invisible until the region exits the staging state.
409                region.switch_state_to_staging(RegionLeaderState::Editing);
410            }
411
412            false
413        } else {
414            let need_compaction = self.config.schedule_compaction_after_edit
415                && edit_result.result.is_ok()
416                && !edit_result.edit.files_to_add.is_empty();
417
418            // Only apply the edit if the result is ok and region is not in staging state.
419            if edit_result.result.is_ok() {
420                // Applies the edit to the region.
421                region.version_control.apply_edit(
422                    Some(edit_result.edit),
423                    &[],
424                    region.file_purger.clone(),
425                );
426            }
427            if edit_result.update_region_state {
428                region.switch_state_to_writable(RegionLeaderState::Editing);
429            }
430
431            need_compaction
432        };
433
434        edit_result
435            .waiters
436            .reply_with(|| match &edit_result.result {
437                Ok(()) => Ok(()),
438                Err(e) => Err(e.clone()).context(EditRegionSnafu {
439                    region_id: edit_result.region_id,
440                }),
441            });
442
443        if edit_result.update_region_state {
444            // Writes stalled specifically by this edit are handled before the next queued edit.
445            // Otherwise the next edit could reserve a committed sequence before those writes.
446            self.handle_region_stalled_requests(&edit_result.region_id, false)
447                .await;
448        }
449
450        let next_request =
451            if let Some(edit_queue) = self.region_edit_queues.get_mut(&edit_result.region_id) {
452                let request = edit_queue.dequeue();
453                if edit_queue.is_empty() {
454                    self.region_edit_queues.remove(&edit_result.region_id);
455                }
456                request
457            } else {
458                None
459            };
460        if let Some(request) = next_request {
461            self.handle_region_edit(request);
462        }
463
464        if need_compaction {
465            self.schedule_compaction(&region).await;
466        }
467    }
468
469    /// Writes truncate action to the manifest and then applies it to the region in background.
470    pub(crate) fn handle_manifest_truncate_action(
471        &self,
472        region: MitoRegionRef,
473        truncate: RegionTruncate,
474        sender: OptionOutputTx,
475    ) {
476        // Marks the region as truncating.
477        // This prevents the region from being accessed by other write requests.
478        if let Err(e) = region.set_truncating() {
479            sender.send(Err(e));
480            return;
481        }
482        // Now the region is in truncating state.
483
484        let request_sender = self.sender.clone();
485        let manifest_ctx = region.manifest_ctx.clone();
486        let is_staging = region.is_staging();
487
488        // Updates manifest in background.
489        common_runtime::spawn_global(async move {
490            // Write region truncated to manifest.
491            let action_list =
492                RegionMetaActionList::with_action(RegionMetaAction::Truncate(truncate.clone()));
493
494            let result = manifest_ctx
495                .update_manifest(RegionLeaderState::Truncating, action_list, is_staging)
496                .await
497                .map(|_| ());
498
499            // Sends the result back to the request sender.
500            let truncate_result = TruncateResult {
501                region_id: truncate.region_id,
502                sender,
503                result,
504                kind: truncate.kind,
505            };
506            let _ = request_sender
507                .send(WorkerRequestWithTime::new(WorkerRequest::Background {
508                    region_id: truncate.region_id,
509                    notify: BackgroundNotify::Truncate(truncate_result),
510                }))
511                .await
512                .inspect_err(|_| warn!("failed to send truncate result"));
513        });
514    }
515
516    /// Advances the durable replay frontier before discarding a region's memtables.
517    pub(crate) fn handle_manifest_discard_unflushed_action(
518        &self,
519        region: MitoRegionRef,
520        discarded_entry_id: EntryId,
521        discarded_sequence: SequenceNumber,
522        discarded_rows: u64,
523        discarded_bytes: u64,
524        sender: OptionOutputTx,
525    ) {
526        if let Err(e) = region.set_truncating() {
527            sender.send(Err(e));
528            return;
529        }
530
531        let region_id = region.region_id;
532        let request_sender = self.sender.clone();
533        let manifest_ctx = region.manifest_ctx.clone();
534
535        common_runtime::spawn_global(async move {
536            // The frontier moves to the last written entry and sequence, so replaying the
537            // WAL after a restart skips everything the memtables held.
538            let edit = RegionEdit {
539                files_to_add: Vec::new(),
540                files_to_remove: Vec::new(),
541                timestamp_ms: None,
542                compaction_time_window: None,
543                flushed_entry_id: Some(discarded_entry_id),
544                flushed_sequence: Some(discarded_sequence),
545                committed_sequence: None,
546            };
547            let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
548            let result = manifest_ctx
549                .update_manifest(RegionLeaderState::Truncating, action_list, false)
550                .await
551                .map(|_| ());
552
553            let result = DiscardUnflushedResult {
554                region_id,
555                sender,
556                result,
557                discarded_entry_id,
558                discarded_sequence,
559                discarded_rows,
560                discarded_bytes,
561            };
562            let _ = request_sender
563                .send(WorkerRequestWithTime::new(WorkerRequest::Background {
564                    region_id,
565                    notify: BackgroundNotify::DiscardUnflushed(result),
566                }))
567                .await
568                .inspect_err(|_| warn!("failed to send discard unflushed result"));
569        });
570    }
571
572    /// Writes region change action to the manifest and then applies it to the region in background.
573    pub(crate) fn handle_manifest_region_change(
574        &self,
575        region: MitoRegionRef,
576        change: RegionChange,
577        need_index: bool,
578        new_options: Option<RegionOptions>,
579        sender: OptionOutputTx,
580    ) {
581        // Marks the region as altering.
582        if let Err(e) = region.set_altering() {
583            sender.send(Err(e));
584            return;
585        }
586        let listener = self.listener.clone();
587        let request_sender = self.sender.clone();
588        let is_staging = region.is_staging();
589        // Now the region is in altering state.
590        common_runtime::spawn_global(async move {
591            let new_meta = change.metadata.clone();
592            let action_list = RegionMetaActionList::with_action(RegionMetaAction::Change(change));
593
594            let result = region
595                .manifest_ctx
596                .update_manifest(RegionLeaderState::Altering, action_list, is_staging)
597                .await
598                .map(|_| ());
599            let notify = WorkerRequest::Background {
600                region_id: region.region_id,
601                notify: BackgroundNotify::RegionChange(RegionChangeResult {
602                    region_id: region.region_id,
603                    sender,
604                    result,
605                    new_meta,
606                    need_index,
607                    new_options,
608                }),
609            };
610            listener
611                .on_notify_region_change_result_begin(region.region_id)
612                .await;
613
614            if let Err(res) = request_sender
615                .send(WorkerRequestWithTime::new(notify))
616                .await
617            {
618                warn!(
619                    "Failed to send region change result back to the worker, region_id: {}, res: {:?}",
620                    region.region_id, res
621                );
622            }
623        });
624    }
625
626    fn update_region_version(
627        version_control: &VersionControlRef,
628        new_meta: RegionMetadataRef,
629        new_options: Option<RegionOptions>,
630        memtable_builder_provider: &MemtableBuilderProvider,
631    ) {
632        let options_changed = new_options.is_some();
633        let region_id = new_meta.region_id;
634        if let Some(new_options) = new_options {
635            // Needs to update the region with new format and memtables.
636            // Creates a new memtable builder for the new options as it may change the memtable type.
637            let new_memtable_builder = memtable_builder_provider.builder_for_options(&new_options);
638            version_control.alter_schema_and_format(new_meta, new_options, new_memtable_builder);
639        } else {
640            // Only changes the schema.
641            version_control.alter_schema(new_meta);
642        }
643
644        let version_data = version_control.current();
645        let version = version_data.version;
646        info!(
647            "Region {} is altered, metadata is {:?}, options: {:?}, options_changed: {}",
648            region_id, version.metadata, version.options, options_changed,
649        );
650    }
651}
652
653/// Checks the edit, writes and applies it.
654async fn edit_region(
655    region: &MitoRegionRef,
656    edit: RegionEdit,
657    cache_manager: CacheManagerRef,
658    listener: WorkerListener,
659    is_staging: bool,
660    preload_sst_cache: bool,
661) -> Result<()> {
662    let region_id = region.region_id;
663    if let Some(write_cache) = cache_manager.write_cache()
664        && preload_sst_cache
665    {
666        for file_meta in &edit.files_to_add {
667            let write_cache = write_cache.clone();
668            let layer = region.access_layer.clone();
669            let listener = listener.clone();
670
671            let index_key = IndexKey::new(region_id, file_meta.file_id, FileType::Parquet);
672            let remote_path =
673                location::sst_file_path(layer.table_dir(), file_meta.file_id(), layer.path_type());
674
675            let is_index_exist = file_meta.exists_index();
676            let index_file_size = file_meta.index_file_size();
677
678            let index_file_index_key = IndexKey::new(
679                region_id,
680                file_meta.index_id().file_id.file_id(),
681                FileType::Puffin(file_meta.index_version),
682            );
683            let index_remote_path = location::index_file_path(
684                layer.table_dir(),
685                file_meta.index_id(),
686                layer.path_type(),
687            );
688
689            let file_size = file_meta.file_size;
690            common_runtime::spawn_global(async move {
691                WRITE_CACHE_INFLIGHT_DOWNLOAD.add(1);
692
693                let parquet_cached = write_cache
694                    .download_if_absent(index_key, &remote_path, layer.object_store(), file_size)
695                    .await;
696
697                if parquet_cached.is_ok() {
698                    // Triggers the filling of the parquet metadata cache.
699                    // The parquet file is already downloaded.
700                    let mut cache_metrics = Default::default();
701                    let _ = write_cache
702                        .file_cache()
703                        .get_parquet_meta_data(
704                            index_key,
705                            &mut cache_metrics,
706                            PageIndexPolicy::Optional,
707                        )
708                        .await;
709
710                    if matches!(parquet_cached, Ok(true)) {
711                        listener.on_file_cache_filled(index_key.file_id);
712                    }
713                }
714                if is_index_exist {
715                    // also download puffin file
716                    if let Err(err) = write_cache
717                        .download(
718                            index_file_index_key,
719                            &index_remote_path,
720                            layer.object_store(),
721                            index_file_size,
722                        )
723                        .await
724                    {
725                        common_telemetry::error!(
726                            err; "Failed to download puffin file, region_id: {}, index_file_index_key: {:?}, index_remote_path: {}", region_id, index_file_index_key, index_remote_path
727                        );
728                    }
729                }
730
731                WRITE_CACHE_INFLIGHT_DOWNLOAD.sub(1);
732            });
733        }
734    }
735
736    info!(
737        "Applying {edit:?} to region {}, is_staging: {}",
738        region_id, is_staging
739    );
740
741    let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
742    region
743        .manifest_ctx
744        .update_manifest(RegionLeaderState::Editing, action_list, is_staging)
745        .await
746        .map(|_| ())
747}