Skip to main content

mito2/worker/
handle_copy_region.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
15use common_telemetry::{debug, error, info};
16use snafu::OptionExt;
17use store_api::region_engine::MitoCopyRegionFromResponse;
18use store_api::storage::{FileId, RegionId};
19
20use crate::error::{InvalidRequestSnafu, MissingManifestSnafu, Result};
21use crate::manifest::action::{RegionEdit, RegionMetaAction, RegionMetaActionList};
22use crate::region::{
23    FileDescriptor, MitoRegionRef, RegionFileCopier, RegionLeaderState, RegionMetadataLoader,
24};
25use crate::request::{
26    BackgroundNotify, CopyRegionFromFinished, CopyRegionFromRequest, WorkerRequest,
27};
28use crate::sst::file::FileMeta;
29use crate::sst::location::region_dir_from_table_dir;
30use crate::worker::{RegionWorkerLoop, WorkerRequestWithTime};
31
32impl<S> RegionWorkerLoop<S> {
33    pub(crate) fn handle_copy_region_from_request(&mut self, request: CopyRegionFromRequest) {
34        let region_id = request.region_id;
35        let source_region_id = request.source_region_id;
36        let sender = request.sender;
37        let region = match self.regions.writable_non_staging_region(region_id) {
38            Ok(region) => region,
39            Err(e) => {
40                let _ = sender.send(Err(e));
41                return;
42            }
43        };
44
45        let same_table = source_region_id.table_id() == region_id.table_id();
46        if !same_table {
47            let _ = sender.send(
48                InvalidRequestSnafu {
49                    region_id,
50                    reason: format!("Source and target regions must be from the same table, source_region_id: {source_region_id}, target_region_id: {region_id}"),
51                }
52                .fail(),
53            );
54            return;
55        }
56        if source_region_id == region_id {
57            let _ = sender.send(
58                InvalidRequestSnafu {
59                    region_id,
60                    reason: format!("Source and target regions must be different, source_region_id: {source_region_id}, target_region_id: {region_id}"),
61                }
62                .fail(),
63            );
64            return;
65        }
66
67        let region_metadata_loader =
68            RegionMetadataLoader::new(self.config.clone(), self.object_store_manager.clone());
69        let worker_sender = self.sender.clone();
70
71        common_runtime::spawn_global(async move {
72            let (region_edit, source_file_ids) = match Self::copy_region_from(
73                &region,
74                region_metadata_loader,
75                source_region_id,
76                region_id,
77                request.parallelism.max(1),
78            )
79            .await
80            {
81                Ok(region_files) => region_files,
82                Err(e) => {
83                    let _ = sender.send(Err(e));
84                    return;
85                }
86            };
87
88            match region_edit {
89                Some(region_edit) => {
90                    if let Err(e) = worker_sender
91                        .send(WorkerRequestWithTime::new(WorkerRequest::Background {
92                            region_id,
93                            notify: BackgroundNotify::CopyRegionFromFinished(
94                                CopyRegionFromFinished {
95                                    region_id,
96                                    edit: region_edit,
97                                    sender,
98                                },
99                            ),
100                        }))
101                        .await
102                    {
103                        error!(e; "Failed to send copy region from finished notification to worker, region_id: {}", region_id);
104                    }
105                }
106                None => {
107                    let _ = sender.send(Ok(MitoCopyRegionFromResponse {
108                        copied_file_ids: source_file_ids,
109                    }));
110                }
111            }
112        });
113    }
114
115    pub(crate) fn handle_copy_region_from_finished(&mut self, request: CopyRegionFromFinished) {
116        let region_id = request.region_id;
117        let sender = request.sender;
118        let region = match self.regions.writable_region(region_id) {
119            Ok(region) => region,
120            Err(e) => {
121                let _ = sender.send(Err(e));
122                return;
123            }
124        };
125
126        let copied_file_ids = request
127            .edit
128            .files_to_add
129            .iter()
130            .map(|file_meta| file_meta.file_id)
131            .collect();
132
133        region
134            .version_control
135            .apply_edit(Some(request.edit), &[], region.file_purger.clone());
136
137        let _ = sender.send(Ok(MitoCopyRegionFromResponse { copied_file_ids }));
138    }
139
140    /// Returns the region edit and the file ids that were copied from the source region to the target region.
141    ///
142    /// If no need to copy files, returns (None, source_file_ids).
143    async fn copy_region_from(
144        region: &MitoRegionRef,
145        region_metadata_loader: RegionMetadataLoader,
146        source_region_id: RegionId,
147        target_region_id: RegionId,
148        parallelism: usize,
149    ) -> Result<(Option<RegionEdit>, Vec<FileId>)> {
150        let table_dir = region.table_dir();
151        let path_type = region.path_type();
152        let region_dir = region_dir_from_table_dir(table_dir, source_region_id, path_type);
153        info!(
154            "Loading source region manifest from region dir: {region_dir}, target region: {target_region_id}"
155        );
156        let source_region_manifest = region_metadata_loader
157            .load_manifest(&region_dir, &region.version().options.storage)
158            .await?
159            .context(MissingManifestSnafu {
160                region_id: source_region_id,
161            })?;
162        let mut new_file_metas = vec![];
163        let target_region_manifest = region.manifest_ctx.manifest().await;
164        let source_file_ids = source_region_manifest
165            .files
166            .keys()
167            .cloned()
168            .collect::<Vec<_>>();
169        debug!(
170            "source region files: {:?}, source region id: {}",
171            source_region_manifest.files, source_region_id
172        );
173        for (file_id, file_meta) in &source_region_manifest.files {
174            if !target_region_manifest.files.contains_key(file_id) {
175                new_file_metas.push(remap_copied_file_meta(file_meta, target_region_id));
176            }
177        }
178        if new_file_metas.is_empty() {
179            return Ok((None, source_file_ids));
180        }
181
182        let file_descriptors = new_file_metas
183            .iter()
184            .flat_map(file_descriptors_for_meta)
185            .collect();
186        debug!("File descriptors to copy: {:?}", file_descriptors);
187        let copier = RegionFileCopier::new(region.access_layer());
188        // TODO(weny): ensure the target region is empty.
189        copier
190            .copy_files(
191                source_region_id,
192                target_region_id,
193                file_descriptors,
194                parallelism,
195            )
196            .await?;
197        let edit = RegionEdit {
198            files_to_add: new_file_metas,
199            files_to_remove: vec![],
200            timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
201            compaction_time_window: None,
202            flushed_entry_id: None,
203            flushed_sequence: None,
204            committed_sequence: None,
205        };
206        let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit.clone()));
207        info!(
208            "Applying {:?} to region {target_region_id}, reason: CopyRegionFrom",
209            edit
210        );
211        let version = region
212            .manifest_ctx
213            .update_manifest(RegionLeaderState::Writable, action_list, false)
214            .await?;
215        info!(
216            "Successfully update manifest version to {version}, region: {target_region_id}, reason: CopyRegionFrom"
217        );
218
219        Ok((Some(edit), source_file_ids))
220    }
221}
222
223fn remap_copied_file_meta(file_meta: &FileMeta, target_region_id: RegionId) -> FileMeta {
224    let mut new_file_meta = file_meta.clone();
225    new_file_meta.region_id = target_region_id;
226    // The target region has an independent sequence domain: the physical
227    // per-row sequences in the copied file belong to the source region, so they
228    // must not be trusted for exact sequence-range reads on the target. Clear
229    // the `preserve_row_sequence` marker and source-domain max sequence to fail
230    // closed until the scan provably cannot intersect the copied rows (see
231    // `files_allow_exact_sequence_range`); otherwise an exact request would
232    // replay source-domain rows as target sequences, and an unmarked file with
233    // a stale source-domain `sequence` hint could be silently skipped as
234    // "proven disjoint".
235    new_file_meta.preserve_row_sequence = false;
236    new_file_meta.sequence = None;
237    new_file_meta
238}
239
240fn file_descriptors_for_meta(file_meta: &FileMeta) -> Vec<FileDescriptor> {
241    let data = FileDescriptor::Data {
242        file_id: file_meta.file_id,
243        size: file_meta.file_size,
244    };
245    if file_meta.exists_index() {
246        let region_index_id = file_meta.index_id();
247        vec![
248            data,
249            FileDescriptor::Index {
250                file_id: region_index_id.file_id.file_id(),
251                version: region_index_id.version,
252                size: file_meta.index_file_size(),
253            },
254        ]
255    } else {
256        vec![data]
257    }
258}