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        // For every file added through region edit, we should fill the file sequence
320        for file in &mut edit.files_to_add {
321            file.sequence = NonZeroU64::new(file_sequence);
322        }
323
324        // Allow retrieving `is_staging` before spawn the edit region task.
325        let is_staging = region.is_staging();
326        let expect_state = if is_staging {
327            RegionLeaderState::Staging
328        } else {
329            RegionLeaderState::Writable
330        };
331        // Marks the region as editing.
332        if let Err(e) = region.set_editing(expect_state) {
333            let e = Arc::new(e);
334            waiters.reply_with(|| Err(e.clone()).context(EditRegionSnafu { region_id }));
335            return;
336        }
337
338        let request_sender = self.sender.clone();
339        let cache_manager = self.cache_manager.clone();
340        let listener = self.listener.clone();
341        // Now the region is in editing state.
342        // Updates manifest in background.
343        common_runtime::spawn_global(async move {
344            let result = edit_region(
345                &region,
346                edit.clone(),
347                cache_manager,
348                listener,
349                is_staging,
350                preload_sst_cache,
351            )
352            .await
353            .map_err(Arc::new);
354            let notify = WorkerRequest::Background {
355                region_id,
356                notify: BackgroundNotify::RegionEdit(RegionEditResult {
357                    region_id,
358                    waiters,
359                    edit,
360                    result,
361                    // we always need to restore region state after region edit
362                    update_region_state: true,
363                    is_staging,
364                }),
365            };
366
367            // We don't set state back as the worker loop is already exited.
368            if let Err(res) = request_sender
369                .send(WorkerRequestWithTime::new(notify))
370                .await
371            {
372                warn!(
373                    "Failed to send region edit result back to the worker, region_id: {}, res: {:?}",
374                    region_id, res
375                );
376            }
377        });
378    }
379
380    /// Handles region edit result.
381    pub(crate) async fn handle_region_edit_result(&mut self, edit_result: RegionEditResult) {
382        let region = match self.regions.get_region(edit_result.region_id) {
383            Some(region) => region,
384            None => {
385                // Fail writes stalled behind this edit if the region was removed before the
386                // edit-completion notification reached the worker.
387                self.fail_region_stalled_requests_as_not_found(&edit_result.region_id);
388                self.reject_region_edit_queue_as_not_found(edit_result.region_id);
389
390                edit_result.waiters.reply_with(|| {
391                    RegionNotFoundSnafu {
392                        region_id: edit_result.region_id,
393                    }
394                    .fail()
395                });
396                return;
397            }
398        };
399
400        let need_compaction = if edit_result.is_staging {
401            if edit_result.update_region_state {
402                // For staging regions, edits are not applied immediately,
403                // as they remain invisible until the region exits the staging state.
404                region.switch_state_to_staging(RegionLeaderState::Editing);
405            }
406
407            false
408        } else {
409            let need_compaction = self.config.schedule_compaction_after_edit
410                && edit_result.result.is_ok()
411                && !edit_result.edit.files_to_add.is_empty();
412
413            // Only apply the edit if the result is ok and region is not in staging state.
414            if edit_result.result.is_ok() {
415                // Applies the edit to the region.
416                region.version_control.apply_edit(
417                    Some(edit_result.edit),
418                    &[],
419                    region.file_purger.clone(),
420                );
421            }
422            if edit_result.update_region_state {
423                region.switch_state_to_writable(RegionLeaderState::Editing);
424            }
425
426            need_compaction
427        };
428
429        edit_result
430            .waiters
431            .reply_with(|| match &edit_result.result {
432                Ok(()) => Ok(()),
433                Err(e) => Err(e.clone()).context(EditRegionSnafu {
434                    region_id: edit_result.region_id,
435                }),
436            });
437
438        if edit_result.update_region_state {
439            // Writes stalled specifically by this edit are handled before the next queued edit.
440            // Otherwise the next edit could reserve a committed sequence before those writes.
441            self.handle_region_stalled_requests(&edit_result.region_id, false)
442                .await;
443        }
444
445        let next_request =
446            if let Some(edit_queue) = self.region_edit_queues.get_mut(&edit_result.region_id) {
447                let request = edit_queue.dequeue();
448                if edit_queue.is_empty() {
449                    self.region_edit_queues.remove(&edit_result.region_id);
450                }
451                request
452            } else {
453                None
454            };
455        if let Some(request) = next_request {
456            self.handle_region_edit(request);
457        }
458
459        if need_compaction {
460            self.schedule_compaction(&region).await;
461        }
462    }
463
464    /// Writes truncate action to the manifest and then applies it to the region in background.
465    pub(crate) fn handle_manifest_truncate_action(
466        &self,
467        region: MitoRegionRef,
468        truncate: RegionTruncate,
469        sender: OptionOutputTx,
470    ) {
471        // Marks the region as truncating.
472        // This prevents the region from being accessed by other write requests.
473        if let Err(e) = region.set_truncating() {
474            sender.send(Err(e));
475            return;
476        }
477        // Now the region is in truncating state.
478
479        let request_sender = self.sender.clone();
480        let manifest_ctx = region.manifest_ctx.clone();
481        let is_staging = region.is_staging();
482
483        // Updates manifest in background.
484        common_runtime::spawn_global(async move {
485            // Write region truncated to manifest.
486            let action_list =
487                RegionMetaActionList::with_action(RegionMetaAction::Truncate(truncate.clone()));
488
489            let result = manifest_ctx
490                .update_manifest(RegionLeaderState::Truncating, action_list, is_staging)
491                .await
492                .map(|_| ());
493
494            // Sends the result back to the request sender.
495            let truncate_result = TruncateResult {
496                region_id: truncate.region_id,
497                sender,
498                result,
499                kind: truncate.kind,
500            };
501            let _ = request_sender
502                .send(WorkerRequestWithTime::new(WorkerRequest::Background {
503                    region_id: truncate.region_id,
504                    notify: BackgroundNotify::Truncate(truncate_result),
505                }))
506                .await
507                .inspect_err(|_| warn!("failed to send truncate result"));
508        });
509    }
510
511    /// Advances the durable replay frontier before discarding a region's memtables.
512    pub(crate) fn handle_manifest_discard_unflushed_action(
513        &self,
514        region: MitoRegionRef,
515        discarded_entry_id: EntryId,
516        discarded_sequence: SequenceNumber,
517        discarded_rows: u64,
518        discarded_bytes: u64,
519        sender: OptionOutputTx,
520    ) {
521        if let Err(e) = region.set_truncating() {
522            sender.send(Err(e));
523            return;
524        }
525
526        let region_id = region.region_id;
527        let request_sender = self.sender.clone();
528        let manifest_ctx = region.manifest_ctx.clone();
529
530        common_runtime::spawn_global(async move {
531            // The frontier moves to the last written entry and sequence, so replaying the
532            // WAL after a restart skips everything the memtables held.
533            let edit = RegionEdit {
534                files_to_add: Vec::new(),
535                files_to_remove: Vec::new(),
536                timestamp_ms: None,
537                compaction_time_window: None,
538                flushed_entry_id: Some(discarded_entry_id),
539                flushed_sequence: Some(discarded_sequence),
540                committed_sequence: None,
541            };
542            let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
543            let result = manifest_ctx
544                .update_manifest(RegionLeaderState::Truncating, action_list, false)
545                .await
546                .map(|_| ());
547
548            let result = DiscardUnflushedResult {
549                region_id,
550                sender,
551                result,
552                discarded_entry_id,
553                discarded_sequence,
554                discarded_rows,
555                discarded_bytes,
556            };
557            let _ = request_sender
558                .send(WorkerRequestWithTime::new(WorkerRequest::Background {
559                    region_id,
560                    notify: BackgroundNotify::DiscardUnflushed(result),
561                }))
562                .await
563                .inspect_err(|_| warn!("failed to send discard unflushed result"));
564        });
565    }
566
567    /// Writes region change action to the manifest and then applies it to the region in background.
568    pub(crate) fn handle_manifest_region_change(
569        &self,
570        region: MitoRegionRef,
571        change: RegionChange,
572        need_index: bool,
573        new_options: Option<RegionOptions>,
574        sender: OptionOutputTx,
575    ) {
576        // Marks the region as altering.
577        if let Err(e) = region.set_altering() {
578            sender.send(Err(e));
579            return;
580        }
581        let listener = self.listener.clone();
582        let request_sender = self.sender.clone();
583        let is_staging = region.is_staging();
584        // Now the region is in altering state.
585        common_runtime::spawn_global(async move {
586            let new_meta = change.metadata.clone();
587            let action_list = RegionMetaActionList::with_action(RegionMetaAction::Change(change));
588
589            let result = region
590                .manifest_ctx
591                .update_manifest(RegionLeaderState::Altering, action_list, is_staging)
592                .await
593                .map(|_| ());
594            let notify = WorkerRequest::Background {
595                region_id: region.region_id,
596                notify: BackgroundNotify::RegionChange(RegionChangeResult {
597                    region_id: region.region_id,
598                    sender,
599                    result,
600                    new_meta,
601                    need_index,
602                    new_options,
603                }),
604            };
605            listener
606                .on_notify_region_change_result_begin(region.region_id)
607                .await;
608
609            if let Err(res) = request_sender
610                .send(WorkerRequestWithTime::new(notify))
611                .await
612            {
613                warn!(
614                    "Failed to send region change result back to the worker, region_id: {}, res: {:?}",
615                    region.region_id, res
616                );
617            }
618        });
619    }
620
621    fn update_region_version(
622        version_control: &VersionControlRef,
623        new_meta: RegionMetadataRef,
624        new_options: Option<RegionOptions>,
625        memtable_builder_provider: &MemtableBuilderProvider,
626    ) {
627        let options_changed = new_options.is_some();
628        let region_id = new_meta.region_id;
629        if let Some(new_options) = new_options {
630            // Needs to update the region with new format and memtables.
631            // Creates a new memtable builder for the new options as it may change the memtable type.
632            let new_memtable_builder = memtable_builder_provider.builder_for_options(&new_options);
633            version_control.alter_schema_and_format(new_meta, new_options, new_memtable_builder);
634        } else {
635            // Only changes the schema.
636            version_control.alter_schema(new_meta);
637        }
638
639        let version_data = version_control.current();
640        let version = version_data.version;
641        info!(
642            "Region {} is altered, metadata is {:?}, options: {:?}, options_changed: {}",
643            region_id, version.metadata, version.options, options_changed,
644        );
645    }
646}
647
648/// Checks the edit, writes and applies it.
649async fn edit_region(
650    region: &MitoRegionRef,
651    edit: RegionEdit,
652    cache_manager: CacheManagerRef,
653    listener: WorkerListener,
654    is_staging: bool,
655    preload_sst_cache: bool,
656) -> Result<()> {
657    let region_id = region.region_id;
658    if let Some(write_cache) = cache_manager.write_cache()
659        && preload_sst_cache
660    {
661        for file_meta in &edit.files_to_add {
662            let write_cache = write_cache.clone();
663            let layer = region.access_layer.clone();
664            let listener = listener.clone();
665
666            let index_key = IndexKey::new(region_id, file_meta.file_id, FileType::Parquet);
667            let remote_path =
668                location::sst_file_path(layer.table_dir(), file_meta.file_id(), layer.path_type());
669
670            let is_index_exist = file_meta.exists_index();
671            let index_file_size = file_meta.index_file_size();
672
673            let index_file_index_key = IndexKey::new(
674                region_id,
675                file_meta.index_id().file_id.file_id(),
676                FileType::Puffin(file_meta.index_version),
677            );
678            let index_remote_path = location::index_file_path(
679                layer.table_dir(),
680                file_meta.index_id(),
681                layer.path_type(),
682            );
683
684            let file_size = file_meta.file_size;
685            common_runtime::spawn_global(async move {
686                WRITE_CACHE_INFLIGHT_DOWNLOAD.add(1);
687
688                let parquet_cached = write_cache
689                    .download_if_absent(index_key, &remote_path, layer.object_store(), file_size)
690                    .await;
691
692                if parquet_cached.is_ok() {
693                    // Triggers the filling of the parquet metadata cache.
694                    // The parquet file is already downloaded.
695                    let mut cache_metrics = Default::default();
696                    let _ = write_cache
697                        .file_cache()
698                        .get_parquet_meta_data(
699                            index_key,
700                            &mut cache_metrics,
701                            PageIndexPolicy::Optional,
702                        )
703                        .await;
704
705                    if matches!(parquet_cached, Ok(true)) {
706                        listener.on_file_cache_filled(index_key.file_id);
707                    }
708                }
709                if is_index_exist {
710                    // also download puffin file
711                    if let Err(err) = write_cache
712                        .download(
713                            index_file_index_key,
714                            &index_remote_path,
715                            layer.object_store(),
716                            index_file_size,
717                        )
718                        .await
719                    {
720                        common_telemetry::error!(
721                            err; "Failed to download puffin file, region_id: {}, index_file_index_key: {:?}, index_remote_path: {}", region_id, index_file_index_key, index_remote_path
722                        );
723                    }
724                }
725
726                WRITE_CACHE_INFLIGHT_DOWNLOAD.sub(1);
727            });
728        }
729    }
730
731    info!(
732        "Applying {edit:?} to region {}, is_staging: {}",
733        region_id, is_staging
734    );
735
736    let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
737    region
738        .manifest_ctx
739        .update_manifest(RegionLeaderState::Editing, action_list, is_staging)
740        .await
741        .map(|_| ())
742}