1use std::sync::Arc;
18use std::time::Instant;
19
20use common_telemetry::{info, warn};
21use object_store::util::{join_path, normalize_dir};
22use snafu::{OptionExt, ResultExt};
23use store_api::logstore::LogStore;
24use store_api::region_request::{AffectedRows, RegionCleanUpRequest, RegionOpenRequest};
25use store_api::storage::RegionId;
26use table::requests::STORAGE_KEY;
27
28use crate::access_layer::AccessLayer;
29use crate::engine::region_hook::{RegionGcInfo, RegionHookRef};
30use crate::error::{
31 ObjectStoreNotFoundSnafu, OpenDalSnafu, OpenRegionSnafu, RegionBusySnafu, RegionNotFoundSnafu,
32 Result,
33};
34use crate::region::opener::{
35 RegionOpener, get_object_store, provider_from_wal_options, sanitize_open_request_options,
36};
37use crate::region::options::RegionOptions;
38use crate::request::OptionOutputTx;
39use crate::sst::location::region_dir_from_table_dir;
40use crate::wal::entry_distributor::WalEntryReceiver;
41use crate::worker::handle_drop::{
42 cleanup_region_file_artifacts, remove_region_dir_for_full_drop, remove_region_dir_once,
43};
44use crate::worker::{DROPPING_MARKER_FILE, RegionWorkerLoop};
45
46impl<S: LogStore> RegionWorkerLoop<S> {
47 pub(crate) async fn handle_offline_cleanup_request(
48 &mut self,
49 region_id: RegionId,
50 mut request: RegionCleanUpRequest,
51 ) -> Result<AffectedRows> {
52 info!(
53 "Try to clean region {} offline, worker: {}",
54 region_id, self.id
55 );
56
57 if self.regions.is_region_exists(region_id) {
58 return RegionBusySnafu { region_id }.fail();
59 }
60
61 sanitize_open_request_options(&mut request.options);
62
63 let options = RegionOptions::try_from_options(region_id, &request.options)?;
64 let object_store = get_object_store(&options.storage, &self.object_store_manager)?;
65 let provider = provider_from_wal_options::<S>(region_id, &options.wal_options)?;
66 self.wal.obsolete_all(region_id, &provider).await?;
67
68 let table_dir = normalize_dir(&request.table_dir);
69 let region_dir = region_dir_from_table_dir(&table_dir, region_id, request.path_type);
70 remove_region_dir_for_full_drop(®ion_dir, &object_store).await?;
71
72 if let Some(hook) = self.plugins.get::<RegionHookRef>() {
89 let access_layer = Arc::new(AccessLayer::new(
90 table_dir.clone(),
91 request.path_type,
92 object_store.clone(),
93 self.puffin_manager_factory.clone(),
94 self.intermediate_manager.clone(),
95 ));
96 let gc_info = RegionGcInfo {
97 removed_files: &[],
98 is_region_dropped: true,
99 full_file_listing: true,
100 };
101 if let Err(err) = hook
102 .on_region_gc(region_id, None, &access_layer, &gc_info)
103 .await
104 {
105 warn!(
106 err;
107 "Region hook on_region_gc failed during offline cleanup for region {}; \
108 returning the error so the caller retries the cleanup",
109 region_id,
110 );
111 return Err(err);
112 }
113 }
114
115 self.cleanup_dropped_region_runtime_state(region_id).await;
116 self.dropping_regions.remove_region(region_id);
117 cleanup_region_file_artifacts(
118 region_id,
119 &table_dir,
120 &self.intermediate_manager,
121 &self.cache_manager,
122 )
123 .await;
124
125 Ok(0)
126 }
127
128 async fn check_and_cleanup_region(
129 &self,
130 region_id: RegionId,
131 request: &RegionOpenRequest,
132 ) -> Result<()> {
133 let object_store = if let Some(storage_name) = request.options.get(STORAGE_KEY) {
134 self.object_store_manager
135 .find(storage_name)
136 .with_context(|| ObjectStoreNotFoundSnafu {
137 object_store: storage_name.clone(),
138 })?
139 } else {
140 self.object_store_manager.default_object_store()
141 };
142 let region_dir =
144 region_dir_from_table_dir(&request.table_dir, region_id, request.path_type);
145 if !self.dropping_regions.is_region_exists(region_id)
146 && object_store
147 .exists(&join_path(®ion_dir, DROPPING_MARKER_FILE))
148 .await
149 .context(OpenDalSnafu)?
150 {
151 let result = remove_region_dir_once(®ion_dir, object_store, true).await;
152 info!(
153 "Region {} is dropped, worker: {}, result: {:?}",
154 region_id, self.id, result
155 );
156 return RegionNotFoundSnafu { region_id }.fail();
157 }
158
159 Ok(())
160 }
161
162 pub(crate) async fn handle_open_request(
163 &mut self,
164 region_id: RegionId,
165 mut request: RegionOpenRequest,
166 wal_entry_receiver: Option<WalEntryReceiver>,
167 sender: OptionOutputTx,
168 ) {
169 if self.regions.is_region_exists(region_id) {
170 sender.send(Ok(0));
171 return;
172 }
173 let Some(sender) = self
174 .opening_regions
175 .wait_for_opening_region(region_id, sender)
176 else {
177 return;
178 };
179 info!("Try to open region {}, worker: {}", region_id, self.id);
180 sanitize_open_request_options(&mut request.options);
181
182 let requirements = request.requirements;
184 let opener = match RegionOpener::new(
185 region_id,
186 &request.table_dir,
187 request.path_type,
188 self.memtable_builder_provider.clone(),
189 self.object_store_manager.clone(),
190 self.purge_scheduler.clone(),
191 self.puffin_manager_factory.clone(),
192 self.intermediate_manager.clone(),
193 self.time_provider.clone(),
194 self.file_ref_manager.clone(),
195 self.partition_expr_fetcher.clone(),
196 )
197 .skip_wal_replay(request.skip_wal_replay)
198 .cache(Some(self.cache_manager.clone()))
199 .hook(self.plugins.get())
200 .wal_entry_reader(wal_entry_receiver.map(|receiver| Box::new(receiver) as _))
201 .replay_checkpoint(request.checkpoint.map(|checkpoint| checkpoint.entry_id))
202 .parse_options(request.options.clone())
203 {
204 Ok(opener) => opener,
205 Err(err) => {
206 sender.send(Err(err));
207 return;
208 }
209 };
210
211 if let Err(err) = opener.ensure_region_requirements(requirements) {
212 sender.send(Err(err));
213 return;
214 }
215
216 if let Err(err) = self.check_and_cleanup_region(region_id, &request).await {
217 sender.send(Err(err));
218 return;
219 }
220
221 let now = Instant::now();
222 let regions = self.regions.clone();
223 let wal = self.wal.clone();
224 let config = self.config.clone();
225 let opening_regions = self.opening_regions.clone();
226 let region_count = self.region_count.clone();
227 let worker_id = self.id;
228 opening_regions.insert_sender(region_id, sender);
229 common_runtime::spawn_global(async move {
230 match opener.open(&config, &wal).await {
231 Ok(region) => {
232 info!(
233 "Region {} is opened with requirements {:?}, worker: {}, elapsed: {:?}",
234 region_id,
235 requirements,
236 worker_id,
237 now.elapsed()
238 );
239 region_count.inc();
240
241 if let Some(hook) = region.manifest_ctx.hook() {
245 hook.on_region_opened(region_id, ®ion.metadata()).await;
246 }
247
248 regions.insert_region(region);
250
251 let senders = opening_regions.remove_sender(region_id);
252 for sender in senders {
253 sender.send(Ok(0));
254 }
255 }
256 Err(err) => {
257 let senders = opening_regions.remove_sender(region_id);
258 let err = Arc::new(err);
259 for sender in senders {
260 sender.send(Err(err.clone()).context(OpenRegionSnafu));
261 }
262 }
263 }
264 });
265 }
266}