1use std::sync::Arc;
16use std::time::Instant;
17
18use common_base::readable_size::ReadableSize;
19use common_telemetry::{debug, error, info};
20use futures::future::try_join_all;
21use object_store::manager::ObjectStoreManagerRef;
22use snafu::{ResultExt, ensure};
23use store_api::metadata::RegionMetadataRef;
24use store_api::region_request::PathType;
25use store_api::storage::{FileId, IndexVersion, RegionId};
26
27use crate::access_layer::AccessLayerRef;
28use crate::config::MitoConfig;
29use crate::error::{self, InvalidSourceAndTargetRegionSnafu, Result};
30use crate::manifest::action::RegionManifest;
31use crate::manifest::manager::{RegionManifestManager, RegionManifestOptions};
32use crate::region::opener::get_object_store;
33use crate::region::options::RegionOptions;
34use crate::sst::file::{RegionFileId, RegionIndexId};
35use crate::sst::location;
36
37#[derive(Debug, Clone)]
39pub struct RegionMetadataLoader {
40 config: Arc<MitoConfig>,
41 object_store_manager: ObjectStoreManagerRef,
42}
43
44impl RegionMetadataLoader {
45 pub fn new(config: Arc<MitoConfig>, object_store_manager: ObjectStoreManagerRef) -> Self {
47 Self {
48 config,
49 object_store_manager,
50 }
51 }
52
53 pub async fn load(
55 &self,
56 region_dir: &str,
57 region_options: &RegionOptions,
58 ) -> Result<Option<RegionMetadataRef>> {
59 let manifest = self
60 .load_manifest(region_dir, ®ion_options.storage)
61 .await?;
62 Ok(manifest.map(|m| m.metadata.clone()))
63 }
64
65 pub async fn load_manifest(
67 &self,
68 region_dir: &str,
69 storage: &Option<String>,
70 ) -> Result<Option<Arc<RegionManifest>>> {
71 let object_store = get_object_store(storage, &self.object_store_manager)?;
72 let region_manifest_options =
73 RegionManifestOptions::new(&self.config, region_dir, &object_store);
74 let Some(manifest_manager) =
75 RegionManifestManager::open(region_manifest_options, &Default::default()).await?
76 else {
77 return Ok(None);
78 };
79
80 let manifest = manifest_manager.manifest();
81 Ok(Some(manifest))
82 }
83}
84
85#[derive(Debug, Clone)]
87pub struct RegionFileCopier {
88 access_layer: AccessLayerRef,
89}
90
91#[derive(Debug, Clone, Copy)]
93pub enum FileDescriptor {
94 Index {
96 file_id: FileId,
97 version: IndexVersion,
98 size: u64,
99 },
100 Data { file_id: FileId, size: u64 },
102}
103
104impl FileDescriptor {
105 pub fn size(&self) -> u64 {
106 match self {
107 FileDescriptor::Index { size, .. } => *size,
108 FileDescriptor::Data { size, .. } => *size,
109 }
110 }
111
112 fn file_path(&self, region_id: RegionId, table_dir: &str, path_type: PathType) -> String {
113 match *self {
114 FileDescriptor::Index {
115 file_id, version, ..
116 } => location::index_file_path(
117 table_dir,
118 RegionIndexId::new(RegionFileId::new(region_id, file_id), version),
119 path_type,
120 ),
121 FileDescriptor::Data { file_id, .. } => {
122 location::sst_file_path(table_dir, RegionFileId::new(region_id, file_id), path_type)
123 }
124 }
125 }
126}
127
128fn build_copy_file_paths(
140 source_region_id: RegionId,
141 target_region_id: RegionId,
142 file_descriptor: FileDescriptor,
143 table_dir: &str,
144 path_type: PathType,
145) -> (String, String) {
146 (
147 file_descriptor.file_path(source_region_id, table_dir, path_type),
148 file_descriptor.file_path(target_region_id, table_dir, path_type),
149 )
150}
151
152fn build_delete_file_path(
153 target_region_id: RegionId,
154 file_descriptor: FileDescriptor,
155 table_dir: &str,
156 path_type: PathType,
157) -> String {
158 file_descriptor.file_path(target_region_id, table_dir, path_type)
159}
160
161impl RegionFileCopier {
162 pub fn new(access_layer: AccessLayerRef) -> Self {
163 Self { access_layer }
164 }
165
166 pub async fn copy_files(
174 &self,
175 source_region_id: RegionId,
176 target_region_id: RegionId,
177 file_ids: Vec<FileDescriptor>,
178 parallelism: usize,
179 ) -> Result<()> {
180 ensure!(
181 source_region_id.table_id() == target_region_id.table_id(),
182 InvalidSourceAndTargetRegionSnafu {
183 source_region_id,
184 target_region_id,
185 },
186 );
187 let table_dir = self.access_layer.table_dir();
188 let path_type = self.access_layer.path_type();
189 let object_store = self.access_layer.object_store();
190
191 info!(
192 "Copying {} files from region {} to region {}",
193 file_ids.len(),
194 source_region_id,
195 target_region_id
196 );
197 debug!(
198 "Copying files: {:?} from region {} to region {}",
199 file_ids, source_region_id, target_region_id
200 );
201 let mut tasks = Vec::with_capacity(parallelism);
202 for skip in 0..parallelism {
203 let target_file_ids = file_ids.iter().skip(skip).step_by(parallelism).copied();
204 let object_store = object_store.clone();
205 tasks.push(async move {
206 for file_desc in target_file_ids {
207 let (source_path, target_path) = build_copy_file_paths(
208 source_region_id,
209 target_region_id,
210 file_desc,
211 table_dir,
212 path_type,
213 );
214 let now = Instant::now();
215 object_store
216 .copy(&source_path, &target_path)
217 .await
218 .inspect_err(
219 |e| error!(e; "Failed to copy file {} to {}", source_path, target_path),
220 )
221 .context(error::OpenDalSnafu)?;
222 let file_size = ReadableSize(file_desc.size());
223 info!(
224 "Copied file {} to {}, file size: {}, elapsed: {:?}",
225 source_path,
226 target_path,
227 file_size,
228 now.elapsed(),
229 );
230 }
231
232 Ok(())
233 });
234 }
235
236 if let Err(err) = try_join_all(tasks).await {
237 error!(err; "Failed to copy files from region {} to region {}", source_region_id, target_region_id);
238 self.clean_target_region(target_region_id, file_ids).await;
239 return Err(err);
240 }
241
242 Ok(())
243 }
244
245 async fn clean_target_region(&self, target_region_id: RegionId, file_ids: Vec<FileDescriptor>) {
247 let table_dir = self.access_layer.table_dir();
248 let path_type = self.access_layer.path_type();
249 let object_store = self.access_layer.object_store();
250 let delete_file_path = file_ids
251 .into_iter()
252 .map(|file_descriptor| {
253 build_delete_file_path(target_region_id, file_descriptor, table_dir, path_type)
254 })
255 .collect::<Vec<_>>();
256 debug!(
257 "Deleting files: {:?} after failed to copy files to target region {}",
258 delete_file_path, target_region_id
259 );
260 if let Err(err) = object_store.delete_iter(delete_file_path).await {
261 error!(err; "Failed to delete files from region {}", target_region_id);
262 }
263 }
264}
265
266#[cfg(test)]
267mod tests {
268 #[cfg(feature = "hdfs-object-store")]
269 use object_store::ObjectStore;
270 #[cfg(feature = "hdfs-object-store")]
271 use object_store::layers::HdfsCompatibilityLayer;
272 #[cfg(feature = "hdfs-object-store")]
273 use object_store::services::Fs;
274
275 use super::*;
276 #[cfg(feature = "hdfs-object-store")]
277 use crate::access_layer::AccessLayer;
278 #[cfg(feature = "hdfs-object-store")]
279 use crate::sst::index::intermediate::IntermediateManager;
280 #[cfg(feature = "hdfs-object-store")]
281 use crate::sst::index::puffin_manager::PuffinManagerFactory;
282
283 #[test]
284 fn test_build_copy_file_paths() {
285 common_telemetry::init_default_ut_logging();
286 let file_id = FileId::random();
287 let source_region_id = RegionId::new(1, 1);
288 let target_region_id = RegionId::new(1, 2);
289 let file_descriptor = FileDescriptor::Data { file_id, size: 100 };
290 let table_dir = "/table_dir";
291 let path_type = PathType::Bare;
292 let (source_path, target_path) = build_copy_file_paths(
293 source_region_id,
294 target_region_id,
295 file_descriptor,
296 table_dir,
297 path_type,
298 );
299 assert_eq!(
300 source_path,
301 format!("/table_dir/1_0000000001/{}.parquet", file_id)
302 );
303 assert_eq!(
304 target_path,
305 format!("/table_dir/1_0000000002/{}.parquet", file_id)
306 );
307
308 let version = 1;
309 let file_descriptor = FileDescriptor::Index {
310 file_id,
311 version,
312 size: 100,
313 };
314 let (source_path, target_path) = build_copy_file_paths(
315 source_region_id,
316 target_region_id,
317 file_descriptor,
318 table_dir,
319 path_type,
320 );
321 assert_eq!(
322 source_path,
323 format!(
324 "/table_dir/1_0000000001/index/{}.{}.puffin",
325 file_id, version
326 )
327 );
328 assert_eq!(
329 target_path,
330 format!(
331 "/table_dir/1_0000000002/index/{}.{}.puffin",
332 file_id, version
333 )
334 );
335 }
336
337 #[test]
338 fn test_build_delete_file_path() {
339 common_telemetry::init_default_ut_logging();
340 let file_id = FileId::random();
341 let target_region_id = RegionId::new(1, 2);
342 let table_dir = "/table_dir";
343 let path_type = PathType::Bare;
344
345 let file_descriptor = FileDescriptor::Data { file_id, size: 100 };
346 let path = build_delete_file_path(target_region_id, file_descriptor, table_dir, path_type);
347 assert_eq!(path, format!("/table_dir/1_0000000002/{}.parquet", file_id));
348
349 let file_descriptor = FileDescriptor::Index {
350 file_id,
351 version: 1,
352 size: 100,
353 };
354 let path = build_delete_file_path(target_region_id, file_descriptor, table_dir, path_type);
355 assert_eq!(
356 path,
357 format!("/table_dir/1_0000000002/index/{}.1.puffin", file_id)
358 );
359 }
360
361 #[cfg(feature = "hdfs-object-store")]
362 #[tokio::test]
363 async fn test_copy_region_files_with_hdfs_fallback() {
364 let (temp_dir, puffin_manager) =
365 PuffinManagerFactory::new_for_test_async("hdfs-copy-region").await;
366 let intermediate_manager = IntermediateManager::init_fs(temp_dir.path().to_string_lossy())
367 .await
368 .unwrap();
369 let storage_dir = temp_dir.path().join("storage");
370 std::fs::create_dir(&storage_dir).unwrap();
371 let object_store = ObjectStore::new(Fs::default().root(storage_dir.to_str().unwrap()))
372 .unwrap()
373 .layer(HdfsCompatibilityLayer::new_for_test());
374 let access_layer = Arc::new(AccessLayer::new(
375 "table_dir",
376 PathType::Bare,
377 object_store.clone(),
378 puffin_manager,
379 intermediate_manager,
380 ));
381 let copier = RegionFileCopier::new(access_layer);
382 let source_region_id = RegionId::new(1, 1);
383 let target_region_id = RegionId::new(1, 2);
384 let file_id = FileId::random();
385 let descriptor = FileDescriptor::Data { file_id, size: 8 };
386 let (source_path, target_path) = build_copy_file_paths(
387 source_region_id,
388 target_region_id,
389 descriptor,
390 "table_dir",
391 PathType::Bare,
392 );
393 object_store.write(&source_path, "contents").await.unwrap();
394
395 copier
396 .copy_files(source_region_id, target_region_id, vec![descriptor], 1)
397 .await
398 .unwrap();
399
400 assert_eq!(
401 b"contents",
402 object_store
403 .read(&target_path)
404 .await
405 .unwrap()
406 .to_bytes()
407 .as_ref()
408 );
409 }
410}