1use 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 ®ion,
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 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(®ion_dir, ®ion.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 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 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}