1use std::fmt;
16use std::sync::Arc;
17
18use common_telemetry::error;
19use store_api::region_request::PathType;
20use store_api::storage::FileId;
21
22use crate::access_layer::AccessLayerRef;
23use crate::cache::CacheManagerRef;
24use crate::error::Result;
25use crate::schedule::scheduler::SchedulerRef;
26use crate::sst::file::{FileMeta, delete_files, delete_index};
27use crate::sst::file_ref::FileReferenceManagerRef;
28use crate::sst::range_index::RangeIndexDeleter;
29
30pub trait FilePurger: Send + Sync + fmt::Debug {
32 fn remove_file(&self, file_meta: FileMeta, is_delete: bool, index_outdated: bool);
37
38 fn new_file(&self, _: &FileMeta) {
42 }
44}
45
46pub type FilePurgerRef = Arc<dyn FilePurger>;
47
48#[derive(Debug)]
50pub struct NoopFilePurger;
51
52impl FilePurger for NoopFilePurger {
53 fn remove_file(&self, _file_meta: FileMeta, _is_delete: bool, _index_outdated: bool) {
54 }
56}
57
58pub struct LocalFilePurger {
60 scheduler: SchedulerRef,
61 sst_layer: AccessLayerRef,
62 cache_manager: Option<CacheManagerRef>,
63 range_index_deleter: Option<RangeIndexDeleter>,
64}
65
66impl fmt::Debug for LocalFilePurger {
67 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
68 f.debug_struct("LocalFilePurger")
69 .field("sst_layer", &self.sst_layer)
70 .finish()
71 }
72}
73
74#[cfg(not(debug_assertions))]
75pub fn should_enable_gc(global_gc_enabled: bool, object_store_scheme: &'static str) -> bool {
77 global_gc_enabled && object_store_scheme != object_store::services::FS_SCHEME
78}
79
80#[cfg(debug_assertions)]
81pub fn should_enable_gc(global_gc_enabled: bool, _object_store_scheme: &'static str) -> bool {
84 global_gc_enabled
85}
86
87pub fn create_file_purger(
98 gc_enabled: bool,
99 path_type: PathType,
100 scheduler: SchedulerRef,
101 sst_layer: AccessLayerRef,
102 cache_manager: Option<CacheManagerRef>,
103 file_ref_manager: FileReferenceManagerRef,
104 range_index_deleter: Option<RangeIndexDeleter>,
105) -> FilePurgerRef {
106 if should_enable_gc(gc_enabled, sst_layer.object_store().info().scheme())
110 && matches!(path_type, PathType::Data | PathType::Bare)
111 {
112 Arc::new(ObjectStoreFilePurger {
113 file_ref_manager,
114 scheduler,
115 range_index_deleter,
116 })
117 } else {
118 Arc::new(
119 LocalFilePurger::new(scheduler, sst_layer, cache_manager)
120 .with_range_index_deleter(range_index_deleter),
121 )
122 }
123}
124
125pub fn create_local_file_purger(
127 scheduler: SchedulerRef,
128 sst_layer: AccessLayerRef,
129 cache_manager: Option<CacheManagerRef>,
130 _file_ref_manager: FileReferenceManagerRef,
131 range_index_deleter: Option<RangeIndexDeleter>,
132) -> FilePurgerRef {
133 Arc::new(
134 LocalFilePurger::new(scheduler, sst_layer, cache_manager)
135 .with_range_index_deleter(range_index_deleter),
136 )
137}
138
139impl LocalFilePurger {
140 pub fn new(
142 scheduler: SchedulerRef,
143 sst_layer: AccessLayerRef,
144 cache_manager: Option<CacheManagerRef>,
145 ) -> Self {
146 Self {
147 scheduler,
148 sst_layer,
149 cache_manager,
150 range_index_deleter: None,
151 }
152 }
153
154 pub fn with_range_index_deleter(mut self, deleter: Option<RangeIndexDeleter>) -> Self {
156 self.range_index_deleter = deleter;
157 self
158 }
159
160 pub async fn stop_scheduler(&self) -> Result<()> {
162 self.scheduler.stop(true).await
163 }
164
165 fn delete_file(&self, file_meta: FileMeta) {
167 let sst_layer = self.sst_layer.clone();
168 let cache_manager = self.cache_manager.clone();
169 if let Err(e) = self.scheduler.schedule(Box::pin(async move {
170 if let Err(e) = delete_files(
171 file_meta.region_id,
172 &[(file_meta.file_id, file_meta.index_id().version)],
173 file_meta.exists_index(),
174 &sst_layer,
175 &cache_manager,
176 )
177 .await
178 {
179 error!(e; "Failed to delete file {:?} from storage", file_meta);
180 }
181 })) {
182 error!(e; "Failed to schedule the file purge request");
183 }
184 }
185
186 fn delete_index(&self, file_meta: FileMeta) {
187 let sst_layer = self.sst_layer.clone();
188 let cache_manager = self.cache_manager.clone();
189 if let Err(e) = self.scheduler.schedule(Box::pin(async move {
190 let index_id = file_meta.index_id();
191 if let Err(e) = delete_index(index_id, &sst_layer, &cache_manager).await {
192 error!(e; "Failed to delete index for file {:?} from storage", file_meta);
193 }
194 })) {
195 error!(e; "Failed to schedule the index purge request");
196 }
197 }
198}
199
200impl FilePurger for LocalFilePurger {
201 fn remove_file(&self, file_meta: FileMeta, is_delete: bool, index_outdated: bool) {
202 if is_delete {
203 schedule_range_index_deletion(
204 &self.scheduler,
205 self.range_index_deleter.as_ref(),
206 file_meta.file_id,
207 );
208 self.delete_file(file_meta);
209 } else if index_outdated {
210 self.delete_index(file_meta);
211 }
212 }
213}
214
215pub struct ObjectStoreFilePurger {
216 file_ref_manager: FileReferenceManagerRef,
217 scheduler: SchedulerRef,
218 range_index_deleter: Option<RangeIndexDeleter>,
219}
220
221impl fmt::Debug for ObjectStoreFilePurger {
222 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
223 f.debug_struct("ObjectStoreFilePurger")
224 .field("file_ref_manager", &self.file_ref_manager)
225 .field("range_index_deleter", &self.range_index_deleter)
226 .finish_non_exhaustive()
227 }
228}
229
230fn schedule_range_index_deletion(
232 scheduler: &SchedulerRef,
233 deleter: Option<&RangeIndexDeleter>,
234 file_id: FileId,
235) {
236 let Some(deleter) = deleter.cloned() else {
237 return;
238 };
239 if let Err(error) = scheduler.schedule(Box::pin(async move {
240 if let Err(error) = deleter.delete(file_id).await {
241 error!(error; "Failed to delete range index, file_id: {file_id}");
242 }
243 })) {
244 error!(error; "Failed to schedule range-index deletion, file_id: {file_id}");
245 }
246}
247
248impl FilePurger for ObjectStoreFilePurger {
249 fn remove_file(&self, file_meta: FileMeta, is_delete: bool, _index_outdated: bool) {
250 self.file_ref_manager.remove_file(&file_meta);
255 if is_delete {
256 schedule_range_index_deletion(
257 &self.scheduler,
258 self.range_index_deleter.as_ref(),
259 file_meta.file_id,
260 );
261 }
262 }
263
264 fn new_file(&self, file_meta: &FileMeta) {
265 self.file_ref_manager.add_file(file_meta);
266 }
267}
268
269#[cfg(test)]
270mod tests {
271 use std::num::NonZeroU64;
272
273 use common_test_util::temp_dir::create_temp_dir;
274 use object_store::ObjectStore;
275 use object_store::services::Fs;
276 use smallvec::SmallVec;
277 use store_api::region_request::PathType;
278 use store_api::storage::{FileId, RegionId};
279
280 use super::*;
281 use crate::access_layer::AccessLayer;
282 use crate::schedule::scheduler::{LocalScheduler, Scheduler};
283 use crate::sst::file::{
284 ColumnIndexMetadata, FileHandle, FileMeta, FileTimeRange, IndexType, RegionFileId,
285 RegionIndexId,
286 };
287 use crate::sst::index::intermediate::IntermediateManager;
288 use crate::sst::index::puffin_manager::PuffinManagerFactory;
289 use crate::sst::location;
290
291 #[tokio::test]
292 async fn test_file_purge() {
293 common_telemetry::init_default_ut_logging();
294
295 let dir = create_temp_dir("file-purge");
296 let dir_path = dir.path().display().to_string();
297 let builder = Fs::default().root(&dir_path);
298 let sst_file_id = RegionFileId::new(RegionId::new(0, 0), FileId::random());
299 let sst_dir = "table1";
300
301 let index_aux_path = dir.path().join("index_aux");
302 let puffin_mgr = PuffinManagerFactory::new(&index_aux_path, 4096, None, None)
303 .await
304 .unwrap();
305 let intm_mgr = IntermediateManager::init_fs(index_aux_path.to_str().unwrap())
306 .await
307 .unwrap();
308
309 let object_store = ObjectStore::new(builder).unwrap();
310
311 let layer = Arc::new(AccessLayer::new(
312 sst_dir,
313 PathType::Bare,
314 object_store.clone(),
315 puffin_mgr,
316 intm_mgr,
317 ));
318 let path = location::sst_file_path(sst_dir, sst_file_id, layer.path_type());
319 object_store.write(&path, vec![0; 4096]).await.unwrap();
320
321 let scheduler = Arc::new(LocalScheduler::new(3));
322
323 let file_purger = Arc::new(LocalFilePurger::new(scheduler.clone(), layer, None));
324
325 {
326 let handle = FileHandle::new(
327 FileMeta {
328 region_id: sst_file_id.region_id(),
329 file_id: sst_file_id.file_id(),
330 time_range: FileTimeRange::default(),
331 level: 0,
332 file_size: 4096,
333 max_row_group_uncompressed_size: 4096,
334 available_indexes: Default::default(),
335 indexes: Default::default(),
336 index_file_size: 0,
337 index_version: 0,
338 num_rows: 0,
339 num_row_groups: 0,
340 sequence: None,
341 partition_expr: None,
342 num_series: 0,
343 ..Default::default()
344 },
345 file_purger,
346 );
347 handle.mark_deleted();
349 }
350
351 scheduler.stop(true).await.unwrap();
352
353 assert!(!object_store.exists(&path).await.unwrap());
354 }
355
356 #[tokio::test]
357 async fn test_file_purge_with_index() {
358 common_telemetry::init_default_ut_logging();
359
360 let dir = create_temp_dir("file-purge");
361 let dir_path = dir.path().display().to_string();
362 let builder = Fs::default().root(&dir_path);
363 let sst_file_id = RegionFileId::new(RegionId::new(0, 0), FileId::random());
364 let index_file_id = RegionIndexId::new(sst_file_id, 0);
365 let sst_dir = "table1";
366
367 let index_aux_path = dir.path().join("index_aux");
368 let puffin_mgr = PuffinManagerFactory::new(&index_aux_path, 4096, None, None)
369 .await
370 .unwrap();
371 let intm_mgr = IntermediateManager::init_fs(index_aux_path.to_str().unwrap())
372 .await
373 .unwrap();
374
375 let object_store = ObjectStore::new(builder).unwrap();
376
377 let layer = Arc::new(AccessLayer::new(
378 sst_dir,
379 PathType::Bare,
380 object_store.clone(),
381 puffin_mgr,
382 intm_mgr,
383 ));
384 let path = location::sst_file_path(sst_dir, sst_file_id, layer.path_type());
385 object_store.write(&path, vec![0; 4096]).await.unwrap();
386
387 let index_path = location::index_file_path(sst_dir, index_file_id, layer.path_type());
388 object_store
389 .write(&index_path, vec![0; 4096])
390 .await
391 .unwrap();
392
393 let scheduler = Arc::new(LocalScheduler::new(3));
394
395 let file_purger = Arc::new(LocalFilePurger::new(scheduler.clone(), layer, None));
396
397 {
398 let handle = FileHandle::new(
399 FileMeta {
400 region_id: sst_file_id.region_id(),
401 file_id: sst_file_id.file_id(),
402 time_range: FileTimeRange::default(),
403 level: 0,
404 file_size: 4096,
405 max_row_group_uncompressed_size: 4096,
406 available_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
407 indexes: vec![ColumnIndexMetadata {
408 column_id: 0,
409 created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
410 }],
411 index_file_size: 4096,
412 index_version: 0,
413 num_rows: 1024,
414 num_row_groups: 1,
415 sequence: NonZeroU64::new(4096),
416 partition_expr: None,
417 num_series: 0,
418 ..Default::default()
419 },
420 file_purger,
421 );
422 handle.mark_deleted();
424 }
425
426 scheduler.stop(true).await.unwrap();
427
428 assert!(!object_store.exists(&path).await.unwrap());
429 assert!(!object_store.exists(&index_path).await.unwrap());
430 }
431}